🙆

AWS Athenaについて

に公開

始めに

業務でAthenaを使う機会が多々ありますが、AIにテーブルなりクエリを作らせてなんとなくで使っている感が強いので整理。

テーブル作成時に最低限必要な情報

テーブル作成時には、大きく以下の情報が必要です。(IAM権限などテーブル作成に直接関係ない箇所は省きます)
①S3の情報
②データファイルの形式
③データスキーマ情報
④CSVファイル特有の情報
⑤パフォーマンス最適化情報

一つずつ深堀していきます。

①S3の情報

データソースとなるS3バケット名、パスは当然必要です。
そして、コスト削減やクエリ時間削減のために、パーティション構造もしっかり確認しておくべきです。
ALBやWAF等のログはパーティションがデフォルトで決まっています。
もし、パーティション構造を変える必要がある場合、Kinesis Data Firehoseを使う必要があります。

②データファイルの形式

ファイル形式(CSV,JSON,Parquet,ORCなど)や圧縮形式、文字エンコーディングも要確認です。
ファイル形式を誤って指定し、テーブルを作成するとAthenaは正しくファイルを読み取ることができなくなります。
また、圧縮形式を指定しないと読み込み性能が劣化するので、適切な圧縮形式を指定するのがおすすめです。

③データスキーマ情報

こちらも当然の話ですが、対象データソースのオブジェクトを一つ見てみて、どのようなスキーマ構造になっているか確認しましょう。
データ構造が間違っていると、正しく集計ができません。
また、NULL値の扱いを間違えるとWHERE条件で意図しないフィルタリングが発生してしまうので注意です。

④CSVファイル特有の情報

対象オブジェクトがCSVの場合は、データを正しくフィルタリングできるようにするために、以下も確認しましょう。
①区切り文字が何か
②ヘッダー行の有無
③引用符が"か'か

⑤パフォーマンス最適化情報

特に継続的にそのテーブルからデータを取得する場合、効率的なクエリを実現しコスト最適化を図るために、以下を確認しましょう。
①データ量: ファイルサイズ、レコード数
②更新頻度: データの追加・更新の頻度
③クエリパターン: 良く検索される条件

上記を確認することで、最適なパーティションやテーブル設計に役立ちます。

パーティション投影

前述の通り、クエリコスト、時間削減のためにパーティションは設定すべきですが、手動でパーティションを管理していくのは大変面倒です。
パーティション投影を使うことで、Athenaがパーティション情報を自動生成してくれるようになるのでぜひ設定しておきましょう!
主要な設定項目を表にしてみました。

プロパティ 説明
projection.enabled 投影機能の有効化 true
projection.{column}.type パーティション型 date, integer, enum
projection.{column}.format 日付フォーマット yyyy/MM/dd
projection.{column}.range 値の範囲 2024/01/01,NOW
projection.{column}.interval 間隔 1
projection.{column}.interval.unit 間隔単位 DAYS, HOURS
storage.location.template S3パステンプレート s3://bucket/data/year=${year}/

例の通り、パーティション投影を組むと、以下のようなパス構造であれば、Athenaがパーティションを自動生成してくれるようになります。

s3://<bucket名>/data/2024/01/01
s3://<bucket名>/data/2024/01/02
s3://<bucket名>/data/2024/01/03

実際のクエリ例(パーティション部抜粋)も記載しておきます。

-- 対JSONファイルの場合を想定
PARTITIONED BY (date string)
ROW FORMAT SERDE 'org.openx.data.jsonserde.JsonSerDe'
LOCATION 's3://my-bucket/data/'
TBLPROPERTIES (
  'projection.enabled'='true',
  'projection.date.type'='date',
  'projection.date.format'='yyyy/MM/dd',
  'projection.date.range'='2024/01/01,NOW',
  'projection.date.interval'='1',
  'projection.date.interval.unit'='DAYS',
  'storage.location.template'='s3://my-bucket/data/date=${date}/'
)

データ加工

今まで記載してきた内容を確認するだけでは、S3のデータを取得することは可能になりますが、データの加工ができません。 多くの場合、データの加工無しに取得したい情報を取得することはできないため、データ加工についてもしっかり理解しておく必要があります。

struct型

AthenaはApache Hiveベースのようなので、struct型を使えます。(AthenaというよりはApache Hiveの話になっていますが、自分のメモ用)
JSONなどの階層データをSQLで効率的に扱うためのデータ型です。

{
  "terminatingrulematchedetails": {
    "conditiontype": "IP_REPUTATION",
    "sensitivitylevel": "HIGH", 
    "locationstring": "HEADER:User-Agent"
  }
}

上記のようなデータがあった場合、以下のようにテーブル定義を書くと、全体が一つの文字列になってしまいます。。。

CREATE TABLE waf_logs (
  terminatingrulematchedetails string
)

structを使うことによって、個々の要素に対してアクセスができるようになります。

CREATE EXTERNAL TABLE waf_logs (
  terminatingrulematchedetails struct<
    conditiontype: string,
    sensitivitylevel: string,
    locationstring: string
  >
)
ROW FORMAT SERDE 'org.openx.data.jsonserde.JsonSerDe'
STORED AS INPUTFORMAT 'org.apache.hadoop.mapred.TextInputFormat'
OUTPUTFORMAT 'org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat'
LOCATION 's3://your-bucket/waf-logs/'

SerDe

シリアライザー、デシリアライザーの略です。
適切なSerDeを使用することで、対象のファイルとSQLテーブル間の翻訳が行えます。

シリアライザー:データを保存用の形式に変更する。
デシリアライザー:ファイルからデータを読み込んで、プログラムで使える形式に変換する。

※Athenaでは主にデシリアライザーが重要です(読み取り専用のため)

SerDeを指定しない場合を考えてみましょう。
WAFログ例の一部抜粋になりますが、S3オブジェクトの中身は以下のようなJSON形式になっています。

{
  "timestamp": 1111111111,
  "action": "BLOCK",
  "httprequest": {
  	"clientip": "192.168.1.1",
  	"country": "JP",
  }
}

これでは、ファイル全体が一つの文字列として認識されてしまい、本当に取得したい情報のみを取得することができません。
SerDeを使うと、以下のように、取得したいキーのみを取得できる形にデータが加工されます。

-- クエリ
SELECT 
  timestamp,
  action,
  httprequest.clientip as client_ip,
  httprequest.country as country
FROM waf_logs
WHERE action = 'BLOCK'
  AND httprequest.country = 'JP'

-- SerDeによる変換後の結果
timestamp: 1111111111
action: "BLOCK"
httprequest.clientip: "192.168.1.1"

※ファイル内容によって、SerDe形式が異なります。

Json:SQL ROW FORMAT SERDE 'org.openx.data.jsonserde.JsonSerDe'
CSV:SQL ROW FORMAT SERDE 'org.apache.hadoop.hive.serde2.lazy.LazySimpleSerde'
Parquet:SQL STORED AS PARQUET

JSON関数

SerDeでは対応しきれない複雑なJSONデータや、動的なJSON構造を扱うための関数です。
例えば、以下のようにログによって構造が異なるようなケースではSerDeでは対応しきれないので、JSON関数を使うといった感じです。

{
  "type": "A",
  "details": {
    "conditiontype": "IP_REPUTATION"
  }
}

{
  "type": "B",
  "error_details": {
    "error_code": "404"
  }
}
-- 実際のクエリ例
SELECT 
  json_extract_scalar(raw_json, '$.type') as log_type,
  CASE 
    WHEN json_extract_scalar(raw_json, '$.type') = 'A'
    THEN json_extract_scalar(raw_json, '$.details.conditiontype')
    ELSE json_extract_scalar(raw_json, '$.error_details.error_code')
  END as extracted_value
FROM raw_logs

INPUTFORMATとOUTPUTFORMAT

SerDeはデータをどう解釈するか(形式の変換)だったのに対し、INPUTFORMATはファイルをどう読み込むか、OUTPUTFORMATはデータをどう書き出すかを定義します。

-- INPUTFORMATの例
-- ①テキストファイル用(1行ずつ読む)
STORED AS INPUTFORMAT 'org.apache.hadoop.mapred.TextInputFormat'
-- ②Parquetファイル用(列指向で読む)
STORED AS PARQUET
-- ③ORCファイル用(圧縮形式で読む)
STORED AS ORC
-- INPUTFORMATの例
-- ①テキスト出力用
STORED AS OUTPUTFORMAT 'org.apache.hive.ql.io.HiveIgnoreKeyTextOutputFormat'
-- ②Parquet出力用
STORED AS OUTPUTFORMAT 'org.apache.hadoop.hive.ql.io.parquet.MapredParquetOutputFormat'

※実際のS3オブジェクトの形式を見て適切なFORMATを指定するようにしましょう。

まとめ

Athenaについて調べてみましたが、結構ややこしいです。。(特に普段DBあまり触っていないのでなおさら)

ひとまず、以下の2つができれば、最低限S3からAthenaでデータを取得することはできるようになると思います。
①調査データの構造、ファイル形式をしっかり確認すること。
②パーティションの切り方をS3パス構造を見て正確に判断すること。

その上で、取得したいデータが素直にテーブルから持ってこれない場合は、データ加工をどうしたらいいのか?頑張って考えましょう!笑

Discussion