😎

Athena と Glue で実現する、とにかくシンプルなデータメッシュの第一歩

に公開

はじめに

データウェアハウスのような中央集権的なアプローチでは、データ管理組織がボトルネックとなる場合があります。「データメッシュ」は、このような組織において特に効果的な手法で、データの所有者とクエリ元が組織や部門ごとに分散しているようなフレームワークといえると思います。

AWS で「データメッシュ」というと、以下の記事にあるような AWS LakeFormation や AWS Data Exchange などが実現のための手段として上がってくるかと思います。

しかし、個人的にはこれらのサービスは冗長なケースも多いと感じています。
例えば、Amazon S3 をデータの境界とし、すでにバケット単位でデータが分類されているようなケースでは、S3 に配置されたデータは共有が前提となるデータです。 S3 に配置するまでのパイプライン上で共有可否を判断しているため、そこに対する AWS LakeFormation や AWS Data Exchange が提供するきめ細やかなアクセス管理手法は冗長に感じられるかもしれません。

データメッシュとは何ですか? - データメッシュアーキテクチャの説明 - AWS
https://aws.amazon.com/jp/what-is/data-mesh/

今回は、そういったユースケースに対応するため、Amazon S3 + Amazon Athena + AWS Glue を活用した簡単なマルチアカウントクエリ構成を具体的に形にしてみました。シンプルな実装例が少ないため、まずは本記事で紹介します。

前提シチュエーション

より具体的なイメージを想起いただくため、以後の設計は以下シチュエーションに対するものとします。

  • 部署は S3 にアップロードする段階で共有すべきデータとそうすべきでないデータを分離している
  • 行レベル、列レベルのアクセス権設定のような細かな権限設定は必要としていない
  • 各部署は、S3 などのサービスについてはある程度知見があるが、追加のサービスについての学習は極力行いたくない、統一的な学習が困難である
  • クエリ元は、各部署のアカウントではなく一旦は特定のアカウントから行うものとする

アーキテクチャ

上記シチュエーションに対して設計したアーキテクチャが以下となります。
このセクションで表示しているアーキテクチャ図は awslabs/diagram-as-code でベースを執筆し、記事用に少し手直ししています。
一つの図に権限関係と、実際に稼働するリソースをすべて入れると見えにくいので、最初に全体図を作った後、確認したい観点ごとに図を生成しています。

全体図

Overall
全体アーキテクチャ図

ポイント

  • サーバーレスの要素を持つ Amazon Athena が Amazon S3 に対するシンプルなクエリエンジンのため採用

IAM 関係以外のリソースのみ示した図


IAM 関係以外のリソースのみ示した図

  • Amazon Athena では、データソース名からフルパスで指定することで、同時に2つ以上のデータソースへの横断的なクエリが可能
  • Amazon Athena でクエリする際は、データの構造情報を提供するカタログとデータの2箇所にアクセスする必要がある点に注意

IAM 関係のリソース図

IAM
IAM 関係のリソース図

  • Amazon Athena は実行ロールで権限の評価を行うため、このロールを起点として各サービスへのアクセス許可が必要
    • Athena のアクセス対象は データの構造情報を提供するカタログ(Glue)とデータ(S3)のため、基本的にはこの2サービスについてアクセス許可が必要
    • AWS Glue はリソースポリシーで、アカウント単位でアクセス元を許可可能
    • Amazon S3 はバケットポリシーでアクセス元の許可が可能

設定までの手順

今回は上記アーキテクチャを実現するサンプルを CDK で作成しました。CDK で構築を行う場合、デプロイするリソースが複数のアカウントにまたがるため少々複雑です。以下に手順を示します。

前半の手順イメージ

  1. 環境変数の設定
export ATHENA_ACCOUNT_ID="アカウントID"
export DATA_ACCOUNT_ID="アカウントID"
export AWS_REGION="ap-northeast-1"
  1. Athena アカウントで CDK Bootstrap
  2. データアカウントで CDK Bootstrap、--trust オプションを使用して、Athena アカウントを指定することで、Athena アカウントから直接 deploy 可能にする

(注意:管理者権限を使用しています)

cdk bootstrap \
  --trust ${ATHENA_ACCOUNT_ID} \
  --cloudformation-execution-policies "arn:aws:iam::aws:policy/AdministratorAccess" \
  aws://${DATA_ACCOUNT_ID}/${AWS_REGION}
  1. Athena アカウントで CDK Deploy
  2. 各アカウントの S3 にデータをアップロード(このステップを飛ばしてもクエリ自体は可能ですが、データがないので結果が0件となります。記事末尾にサンプルを掲載しています)
  3. AWS Console から Amazon Athena を開いてクエリを実行し、各アカウントのデータをクエリして分析できていることを確認する

無事にセットアップできた場合、以下のように Amazon Athena コンソールにデータソースが2つ表示されています。

これを用いて、例えば以下のようにフルパスでデータソースを指定して結合することで、他アカウントのデータを含めたクエリが可能です。

SELECT 
    c.region,
    o.order_id,
    o.total_amount,
    o.order_date
FROM "data_account_catalog".customer_db.customers c
JOIN "athena_account_catalog".order_db.orders o 
    ON c.customer_id = o.customer_id
LIMIT 10;

正しく設定できた場合、以下のような結果が確認できます。

現状の課題と今後

冒頭の記述の通り、今回は既存のソリューションをシンプルにした例を示しました。そのため、より複雑にしたい場合、例えば以下のような要件がある場合は逆に労力が増したりする場合があります。

  1. AWS Glue リソースポリシーは現状 AWS CDK で扱えないのでカスタムリソースが必要
  2. AWS Glue リソースポリシーがアカウント単位での設定になってしまうため、Athena でクエリするアカウントが複数になると、単一のポリシーを重複更新することになる(以下に補足用のの図を記載)
  3. 細かいデータの閲覧制限には向かないので、Amazon S3 に上げた後に制御したいのであれば 冒頭で触れた AWS LakeFormation が望ましいと思われる

執筆後に気づきましたが、AWS Glue リソースポリシーの代わりに Resource Access Manager で Catalogue の所有権を共有する方法もあります。組織内で相互に共有したい場合はこちらがおすすめだろうと思われます。 Resource Access Manager は組織内共有以外に AWS Account を対象に共有することも可能なため、クエリ元が分散する場合にも対応ができそうです。

2に関する補足

Appendix

必要最小限のコードを以下に示します。将来 Github 等で公開した場合はそちらに差し替え予定です。
ある程度権限周りは絞って設定していますが、あくまでサンプルのため実行の際は十分ご注意ください。

app.ts
import 'source-map-support/register';
import * as cdk from 'aws-cdk-lib';
import { AthenaMainStack } from './stacks/athena-main-stack';
import { DataSourceStack } from './stacks/data-source-stack';

const app = new cdk.App();

// アカウントIDを環境変数またはコンテキストから取得(必須)
const athenaAccountId = process.env.ATHENA_ACCOUNT_ID || app.node.tryGetContext('athenaAccountId');
const dataAccountId = process.env.DATA_ACCOUNT_ID || app.node.tryGetContext('dataAccountId');
const region = process.env.AWS_REGION || 'ap-northeast-1';

if (!athenaAccountId || !dataAccountId ) {
  throw new Error('Account ID must be set');
}

// 注文データ(Athenaアカウント)
const ordersStack = new DataSourceStack(app, 'DataSourceStackOrders', {
  env: { account: athenaAccountId, region },
  description: 'Order data source stack (same account as Athena)',
  dataSourceConfig: {
    tableName: 'orders',
    databaseName: 'order_db',
    bucketSuffix: 'orders',
    schema: [
      { name: 'order_id', type: 'string' },
      { name: 'customer_id', type: 'string' },
      { name: 'product_id', type: 'string' },
      { name: 'quantity', type: 'int' },
      { name: 'order_date', type: 'timestamp' },
      { name: 'total_amount', type: 'double' }
    ]
  },
  athenaAccountId
});

// 顧客データ(データアカウント)→ 別アカウントのため自動的にクロスアカウント権限を設定
const customersStack = new DataSourceStack(app, 'DataSourceStackCustomers', {
  env: { account: dataAccountId, region },
  description: 'Customer data source stack',
  dataSourceConfig: {
    tableName: 'customers',
    databaseName: 'customer_db',
    bucketSuffix: 'customers',
    schema: [
      { name: 'customer_id', type: 'string' },
      { name: 'name', type: 'string' },
      { name: 'email', type: 'string' },
      { name: 'region', type: 'string' },
      { name: 'created_date', type: 'timestamp' }
    ]
  },
  athenaAccountId
});


/**
 * AthenaMainStack: Athena WorkGroup、クエリ結果用S3バケット、Athena実行用IAMロールを作成
 * - 全てのデータソーススタックが作成された後に作成される
 */
const athenaMainStack = new AthenaMainStack(app, 'AthenaMainStack', {
  env: { account: athenaAccountId, region },
  description: 'Main Athena stack for cross-account queries with order data',
  dataSourceAccounts: [dataAccountId],
  athenaAccountId: athenaAccountId,
  dataAccountId: dataAccountId
});

// Athena Main Stack は全てのデータソーススタックに依存
athenaMainStack.addDependency(ordersStack);
athenaMainStack.addDependency(customersStack);
athena-main-stack.ts
import * as cdk from 'aws-cdk-lib';
import * as s3 from 'aws-cdk-lib/aws-s3';
import * as athena from 'aws-cdk-lib/aws-athena';
import * as iam from 'aws-cdk-lib/aws-iam';
import { Construct } from 'constructs';

export interface AthenaMainStackProps extends cdk.StackProps {
  dataSourceAccounts: string[];
  athenaAccountId: string;
  dataAccountId: string;
}

/**
 * Athenaクエリを実行するメインアカウントのスタック
 * - Athenaワークグループ
 * - クエリ結果保存用S3バケット
 * - Athenaクエリ実行用IAMロール
 * - クロスアカウントアクセス用のGlue Catalogデータベース
 */
export class AthenaMainStack extends cdk.Stack {
  public readonly queryResultsBucket: s3.Bucket;
  public readonly workGroup: athena.CfnWorkGroup;
  public readonly athenaExecutionRole: iam.Role;

  constructor(scope: Construct, id: string, props: AthenaMainStackProps) {
    super(scope, id, props);

    const { dataSourceAccounts, athenaAccountId, dataAccountId } = props;

    this.queryResultsBucket = new s3.Bucket(this, 'AthenaQueryResultsBucket', {
      bucketName: `athena-query-results-${this.account}-${this.region}`,
      removalPolicy: cdk.RemovalPolicy.DESTROY,
      autoDeleteObjects: true,
      encryption: s3.BucketEncryption.S3_MANAGED,
      versioned: false,
      lifecycleRules: [{
        id: 'delete-after-30-days',
        expiration: cdk.Duration.days(30)
      }]
    });

    // Athenaクエリ実行用IAMロール
    this.athenaExecutionRole = new iam.Role(this, 'AthenaExecutionRole', {
      roleName: `AthenaExecutionRole-${this.region}`,
      assumedBy: new iam.CompositePrincipal(
        new iam.ServicePrincipal('athena.amazonaws.com'),
        new iam.AccountRootPrincipal() // 管理者がAssumeできるようにする
      ),
      description: 'Role for Athena to execute cross-account queries',
      managedPolicies: [
        iam.ManagedPolicy.fromAwsManagedPolicyName('AmazonAthenaFullAccess')
      ]
    });

    const accountsWithSuffixes = [
      ...dataSourceAccounts.map(accountId => ({ accountId, suffixes: ['customers'] })),
      { accountId: this.account, suffixes: ['orders'] }
    ];
    accountsWithSuffixes.forEach(({ accountId, suffixes }) => {
      const s3Resources = suffixes.flatMap(suffix => [
        `arn:aws:s3:::cross-account-data-${suffix}-${accountId}-${this.region}`,
        `arn:aws:s3:::cross-account-data-${suffix}-${accountId}-${this.region}/*`
      ]);

      this.athenaExecutionRole.addToPolicy(new iam.PolicyStatement({
        effect: iam.Effect.ALLOW,
        actions: ['s3:GetObject', 's3:ListBucket', 's3:GetBucketLocation'],
        resources: s3Resources
      }));

      this.athenaExecutionRole.addToPolicy(new iam.PolicyStatement({
        effect: iam.Effect.ALLOW,
        actions: [
          'glue:GetDatabase',
          'glue:GetDatabases',
          'glue:GetTable',
          'glue:GetTables',
          'glue:GetPartition',
          'glue:GetPartitions',
          'glue:BatchCreatePartition',
          'glue:BatchDeletePartition',
          'glue:BatchUpdatePartition'
        ],
        resources: [
          `arn:aws:glue:${this.region}:${accountId}:catalog`,
          `arn:aws:glue:${this.region}:${accountId}:database/*`,
          `arn:aws:glue:${this.region}:${accountId}:table/*/*`
        ]
      }));
    });
    this.queryResultsBucket.grantReadWrite(this.athenaExecutionRole);


    // Athenaワークグループ
    this.workGroup = new athena.CfnWorkGroup(this, 'CrossAccountWorkGroup', {
      name: 'cross-account-queries',
      description: 'Work group for cross-account Athena queries',
      state: 'ENABLED',
      workGroupConfiguration: {
        resultConfiguration: {
          outputLocation: `s3://${this.queryResultsBucket.bucketName}/`,
          encryptionConfiguration: {
            encryptionOption: 'SSE_S3'
          }
        },
        enforceWorkGroupConfiguration: true
      }
    });
    this.workGroup.applyRemovalPolicy(cdk.RemovalPolicy.DESTROY);


    // Athena データソース
    const dataCatalogs = [
      { id: 'AthenaAccountDataSource', name: 'athena_account_catalog', catalogId: athenaAccountId, description: 'Glue Data Catalog in Athena account (same account)' },
      { id: 'DataAccountDataSource', name: 'data_account_catalog', catalogId: dataAccountId, description: 'Glue Data Catalog in Data account (cross-account)' }
    ];

    dataCatalogs.forEach(({ id, name, catalogId, description }) => {
      const catalog = new athena.CfnDataCatalog(this, id, {
        name,
        type: 'GLUE',
        description,
        parameters: { 'catalog-id': catalogId }
      });
      catalog.applyRemovalPolicy(cdk.RemovalPolicy.DESTROY);
    });

  }
}
data-source-stack.ts
import * as cdk from 'aws-cdk-lib';
import * as s3 from 'aws-cdk-lib/aws-s3';
import * as glue from 'aws-cdk-lib/aws-glue';
import * as iam from 'aws-cdk-lib/aws-iam';
import { AwsCustomResource, AwsCustomResourcePolicy, PhysicalResourceId } from 'aws-cdk-lib/custom-resources';
import { Construct } from 'constructs';

export interface ColumnSchema {
  name: string;
  type: string;
}

export interface DataSourceConfig {
  tableName: string;
  databaseName: string;
  bucketSuffix: string;
  schema: ColumnSchema[];
}

export interface DataSourceStackProps extends cdk.StackProps {
  dataSourceConfig: DataSourceConfig;
  athenaAccountId: string;
}


export class DataSourceStack extends cdk.Stack {
  public readonly dataBucket: s3.Bucket;
  public readonly database: glue.CfnDatabase;
  public readonly table: glue.CfnTable;

  constructor(scope: Construct, id: string, props: DataSourceStackProps) {
    super(scope, id, props);
    const { dataSourceConfig, athenaAccountId } = props;

    // データ保存用S3バケット
    this.dataBucket = new s3.Bucket(this, 'DataBucket', {
      bucketName: `cross-account-data-${dataSourceConfig.bucketSuffix}-${this.account}-${this.region}`,
      removalPolicy: cdk.RemovalPolicy.DESTROY,
      autoDeleteObjects: true,
      encryption: s3.BucketEncryption.S3_MANAGED,
    });

    this.database = new glue.CfnDatabase(this, 'Database', {
      catalogId: this.account,
      databaseInput: {
        name: dataSourceConfig.databaseName,
        description: `Database for ${dataSourceConfig.tableName} data`
      }
    });
    this.database.applyRemovalPolicy(cdk.RemovalPolicy.DESTROY);

    this.table = new glue.CfnTable(this, 'Table', {
      catalogId: this.account,
      databaseName: this.database.ref,
      tableInput: {
        name: dataSourceConfig.tableName,
        description: `Table for ${dataSourceConfig.tableName} data`,
        tableType: 'EXTERNAL_TABLE',
        parameters: {
          'classification': 'parquet',
          'compressionType': 'snappy',
          'typeOfData': 'file'
        },
        storageDescriptor: {
          columns: dataSourceConfig.schema.map(column => ({
            name: column.name,
            type: column.type
          })),
          location: `s3://${this.dataBucket.bucketName}/data/`,
          inputFormat: 'org.apache.hadoop.hive.ql.io.parquet.MapredParquetInputFormat',
          outputFormat: 'org.apache.hadoop.hive.ql.io.parquet.MapredParquetOutputFormat',
          serdeInfo: {
            serializationLibrary: 'org.apache.hadoop.hive.ql.io.parquet.serde.ParquetHiveSerDe'
          },
          compressed: true,
          storedAsSubDirectories: false
        }
      }
    });
    this.table.applyRemovalPolicy(cdk.RemovalPolicy.DESTROY);


    // S3バケットにAthenaアカウントからのアクセスを許可
    this.dataBucket.addToResourcePolicy(new iam.PolicyStatement({
      sid: 'AllowCrossAccountAthenaAccess',
      effect: iam.Effect.ALLOW,
      principals: [new iam.AccountPrincipal(athenaAccountId)],
      actions: [
        's3:GetObject',
        's3:ListBucket',
        's3:GetBucketLocation'
      ],
      resources: [
        this.dataBucket.bucketArn,
        `${this.dataBucket.bucketArn}/*`
      ]
    }));

    // Glue にAthenaアカウントからのアクセスを許可
    new AwsCustomResource(this, 'GlueResourcePolicy', {
      onUpdate: {
        service: 'Glue',
        action: 'putResourcePolicy',
        parameters: {
          PolicyInJson: JSON.stringify({
            Version: '2012-10-17',
            Statement: [{
              Sid: 'AllowCrossAccountGlueCatalogAccess',
              Effect: 'Allow',
              Principal: { AWS: `arn:aws:iam::${athenaAccountId}:root` },
              Action: [
                'glue:GetDatabase',
                'glue:GetDatabases',
                'glue:GetTable',
                'glue:GetTables',
                'glue:GetPartition',
                'glue:GetPartitions',
                'glue:BatchCreatePartition',
                'glue:BatchDeletePartition',
                'glue:BatchUpdatePartition'
              ],
              Resource: [
                `arn:aws:glue:${this.region}:${this.account}:catalog`,
                `arn:aws:glue:${this.region}:${this.account}:database/${dataSourceConfig.databaseName}`,
                `arn:aws:glue:${this.region}:${this.account}:table/${dataSourceConfig.databaseName}/${dataSourceConfig.tableName}`
              ]
            }]
          })
        },
        physicalResourceId: PhysicalResourceId.of('glueDataCatalogPermissions')
      },
      onDelete: {
        service: 'Glue',
        action: 'deleteResourcePolicy'
      },
      policy: AwsCustomResourcePolicy.fromSdkCalls({
        resources: AwsCustomResourcePolicy.ANY_RESOURCE
      })
    });

  }
}

generate-sample-data.py
#!/usr/bin/env python3
"""
Cross-Account Athena Query System用のサンプルデータ生成スクリプト
Parquet形式でサンプルデータを生成します。
"""

import pandas as pd
import numpy as np
from datetime import datetime, timedelta
import random
import os

def generate_customers(num_records=1000):
    """顧客データを生成 (架空の名前) """
    fake_names = [
        'alice ', 'bob', 'carol', 'dave',
        'eve', '山本誠', 'frank'
    ]
    
    regions = ['東京', '大阪', '名古屋', '福岡', '札幌', '仙台', '広島', '那覇']
    
    data = []
    for i in range(num_records):
        customer_id = f'CUST_{i+1:06d}'
        name = random.choice(fake_names) + str(random.randint(1, 999))
        email = f'customer{i+1}@example.com'
        region = random.choice(regions)
        created_date = datetime.now() - timedelta(days=random.randint(0, 365*2))
        
        data.append({
            'customer_id': customer_id,
            'name': name,
            'email': email,
            'region': region,
            'created_date': created_date
        })
    
    return pd.DataFrame(data)

def generate_products(num_records=500):
    """商品データを生成"""
    product_names = [
        'ノートパソコン', 'マウス', 'キーボード', 'モニター', 'プリンター',
        'スマートフォン', 'タブレット', 'イヤホン', 'スピーカー', 'Webカメラ',
        'デスクチェア', 'デスクライト', 'ファイル', 'ノート', 'ペン'
    ]
    
    categories = ['電子機器', 'コンピュータ', '周辺機器', 'オフィス用品', 'モバイル']
    suppliers = ['メーカーA', 'メーカーB', 'メーカーC', 'メーカーD', 'メーカーE']
    
    data = []
    for i in range(num_records):
        product_id = f'PROD_{i+1:06d}'
        name = random.choice(product_names) + f' {random.choice(["Pro", "Plus", "Lite", "Standard", "Premium"])}'
        price = round(random.uniform(500, 150000), 2)
        category = random.choice(categories)
        supplier = random.choice(suppliers)
        
        data.append({
            'product_id': product_id,
            'name': name,
            'price': price,
            'category': category,
            'supplier': supplier
        })
    
    return pd.DataFrame(data)

def generate_orders(customers_df, products_df, num_records=5000):
    """注文データを生成"""
    data = []
    customer_ids = customers_df['customer_id'].tolist()
    product_ids = products_df['product_id'].tolist()
    
    for i in range(num_records):
        order_id = f'ORD_{i+1:08d}'
        customer_id = random.choice(customer_ids)
        product_id = random.choice(product_ids)
        quantity = random.randint(1, 10)
        
        # 対応する商品の価格を取得
        product_price = products_df[products_df['product_id'] == product_id]['price'].iloc[0]
        total_amount = round(product_price * quantity, 2)
        
        order_date = datetime.now() - timedelta(days=random.randint(0, 365))
        
        data.append({
            'order_id': order_id,
            'customer_id': customer_id,
            'product_id': product_id,
            'quantity': quantity,
            'order_date': order_date,
            'total_amount': total_amount
        })
    
    return pd.DataFrame(data)

def main():
    """メイン処理"""
    print("Cross-Account Athena Query System - サンプルデータ生成中...")
    
    # スクリプトディレクトリを取得
    script_dir = os.path.dirname(os.path.abspath(__file__))
    project_root = os.path.dirname(script_dir)
    data_dir = os.path.join(project_root, 'data')
    
    # データ出力ディレクトリを作成
    os.makedirs(data_dir, exist_ok=True)
    
    # 顧客データ生成
    print("顧客データを生成中...")
    customers_df = generate_customers(1000)
    customers_df.to_parquet(os.path.join(data_dir, 'customers-sample.parquet'), index=False)
    print(f"顧客データ生成完了: {len(customers_df)}件")
    
    # 商品データ生成
    print("商品データを生成中...")
    products_df = generate_products(500)
    products_df.to_parquet(os.path.join(data_dir, 'products-sample.parquet'), index=False)
    print(f"商品データ生成完了: {len(products_df)}件")
    
    # 注文データ生成
    print("注文データを生成中...")
    orders_df = generate_orders(customers_df, products_df, 5000)
    orders_df.to_parquet(os.path.join(data_dir, 'orders-sample.parquet'), index=False)
    print(f"注文データ生成完了: {len(orders_df)}件")
    
    # データサマリー出力
    print("\n=== データサマリー ===")
    print(f"顧客数: {len(customers_df)}")
    print(f"商品数: {len(products_df)}")
    print(f"注文数: {len(orders_df)}")
    print(f"総売上金額: ¥{orders_df['total_amount'].sum():,.2f}")
    
    # サンプルデータプレビュー
    print("\n=== 顧客データサンプル ===")
    print(customers_df.head())
    
    print("\n=== 商品データサンプル ===")
    print(products_df.head())
    
    print("\n=== 注文データサンプル ===")
    print(orders_df.head())
    
    print("\nサンプルデータ生成が完了しました!")
    print("以下のファイルが生成されました:")
    print(f"- {os.path.join(data_dir, 'customers-sample.parquet')}")
    print(f"- {os.path.join(data_dir, 'products-sample.parquet')}")  
    print(f"- {os.path.join(data_dir, 'orders-sample.parquet')}")

if __name__ == "__main__":
    main()
アマゾン ウェブ サービス ジャパン (有志)

Discussion