AWS ãæŽ»çšããããŒã¿ã¬ã€ã¯ã¯ãããããŠé«ãå¯çšæ§ãèªã Amazon Simple Storage Service(Amazon S3) ãåå°ãšããŠããã倿§ãªããŒã¿ãšã¢ããªãã£ã¯ã¹ã®ã¢ãããŒããçµã¿åãããã®ã«å¿
èŠãªã¹ã±ãŒã«ãææ·æ§ãæè»æ§ãæäŸããããšãã§ããŸããããŒã¿ã¬ã€ã¯ããµã€ãºãšå©çšæ³ã®äž¡é¢ã§æçããŠããã«ã€ããããŒã¿ãããžãã¹ã€ãã³ãã«åãããŠäžè²«æ§ãä¿ã€ã®ã«ããªãã®åŽåãè²»ããããããšããããŸãããã¡ã€ã«ããã©ã³ã¶ã¯ã·ã§ã³ãšäžè²«æ§ãä¿ã£ãп޿°ãããããšã確å®ã«ããããã«ã Apache Iceberg ã Apache Hudi ã Linux Foundation Delta Lake ãªã©ã®ãªãŒãã³ãœãŒã¹ã®ãã©ã³ã¶ã¯ã·ã§ã³åŠçå¯èœãªããŒãã«ãã©ãŒããã (Open Table Format â OTF) ãå©çšãã顧客ãå¢ããŠããŸãããããã®ãã©ãŒãããã¯é«ãå§çž®çã§ããŒã¿ãä¿åããã¢ããªã±ãŒã·ã§ã³ããã¬ãŒã ã¯ãŒã¯ãšé£æºããAmazon S3 äžã«æ§ç¯ãããããŒã¿ã¬ã€ã¯ã§ã®å·®åïŒå¢åïŒããŒã¿åŠçãç°¡çŽ åããŸãããããã®ãã©ãŒãããã«ããã ACID (ååæ§ãäžè²«æ§ãå颿§ãæç¶æ§)ãã©ã³ã¶ã¯ã·ã§ã³ãã¢ãããµãŒããåé€ãã¿ã€ã ãã©ãã«ãã¹ãããã·ã§ãããªã©ã®é«åºŠãªæ©èœãå¯èœã«ãªããŸãããããã®æ©èœã¯ä»¥åã¯ããŒã¿ãŠã§ã¢ããŠã¹ã§ã®ã¿å©çšã§ãããã®ã§ããåããŒãã«ãã©ãŒãããã¯ãã®æ©èœãå°ããã€ç°ãªãæ¹æ³ã§å®è£
ããŠããŸããæ¯èŒã®ããã«ã¯ã AWS äžã®ãã©ã³ã¶ã¯ã·ã§ãã«ããŒã¿ã¬ã€ã¯ã®ããã®ãªãŒãã³ããŒãã«ãã©ãŒãããã®éžæ ïŒè±æïŒãåç
§ããŠãã ããã 2023幎ã«ãAWS 㯠Amazon Athena for Apache Spark ã«ãããŠãApache IcebergãApache HudiãLinux Foundation Delta Lake ãµããŒãã®äžè¬æäŸéå§ã çºè¡š ããŸãããããã«ãããåå¥ã®ã³ãã¯ã¿ãé¢é£ããäŸåé¢ä¿ãã€ã³ã¹ããŒã«ããŠããã±ãŒãžç®¡çããå¿
èŠããªããªãããããã®ãã¬ãŒã ã¯ãŒã¯ã䜿çšããããã«å¿
èŠãªèšå®æé ãç°¡çŽ åãããŸãã ãã®æçš¿ã§ã¯ã Amazon Athena ããŒãããã¯ã§ Spark SQL ã䜿çšããæ¹æ³ãšãIcebergãHudiãDelta Lake ããŒãã«ãã©ãŒããããæäœããæ¹æ³ã瀺ããŸãã Athena ã® Spark SQL ã䜿çšãããããŒã¿ããŒã¹ãšããŒãã«ã®äœæãããŒãã«ãžã®ããŒã¿ã®æ¿å
¥ãããŒã¿ã®ã¯ãšãªãAmazon S3 ã®ããŒãã«ã¹ãããã·ã§ããã®ç¢ºèªãªã©ãäžè¬çãªæäœããã¢ã³ã¹ãã¬ãŒã·ã§ã³ããŸãã åææ¡ä»¶ 次ã®åææ¡ä»¶ãå®äºããŠãã ãã: ã Amazon Athena Spark ã§ Spark SQL ãå®è¡ãã ãïŒè±æïŒã«èšèŒãããŠãããã¹ãŠã®åææ¡ä»¶ãæºãããŠããããšã確èªããŠãã ããã ã Amazon Athena Spark ã§ Spark SQL ãå®è¡ãã ãã§è©³è¿°ãããŠããããã«ãAWS Glue Data Catalog ã« sparkblogdb ãšããããŒã¿ããŒã¹ãš noaa_pq ãšããããŒãã«ãäœæããŠãã ããã Athena ã¯ãŒã¯ã°ã«ãŒãã§äœ¿çšããã AWS Identity and Access Management (IAM) ããŒã«ã«ãS3 ãã±ãããšãã¬ãã£ãã¯ã¹ãžã®èªã¿æžãã¢ã¯ã»ã¹èš±å¯ãä»äžããŠãã ããã 詳现ã«ã€ããŠã¯ã Amazon S3: Allows read and write access to objects in an S3 Bucket ãåç
§ããŠãã ããã ã¯ãªãŒã³ã¢ãããå®è¡ããããã«ãAthena ã¯ãŒã¯ã°ã«ãŒãã§äœ¿çšããã IAM ããŒã«ã«ãS3 ãã±ãããšãã¬ãã£ãã¯ã¹ãžã® s3:DeleteObject ã¢ã¯ã»ã¹èš±å¯ãä»äžããŠãã ããã 詳现ã«ã€ããŠã¯ã Amazon S3 ã¢ã¯ã·ã§ã³ ã® ãªããžã§ã¯ãåé€ã®ã¢ã¯ã»ã¹èš±å¯ ã»ã¯ã·ã§ã³ãåç
§ããŠãã ããã Amazon S3 ããã®ãµã³ãã«ããŒãããã¯ã®ããŠã³ããŒããšã€ã³ããŒã ãã®æçš¿ã§èª¬æããããŒãããã¯ã¯ã次ã®å ŽæããããŠã³ããŒãã§ããŸãã Iceberg ãã¥ãŒããªã¢ã«ããŒãããã¯: s3://athena-examples-us-east-1/athenasparksqlblog/notebooks/SparkSQL_iceberg.ipynb Hudi ãã¥ãŒããªã¢ã«ããŒãããã¯: s3://athena-examples-us-east-1/athenasparksqlblog/notebooks/SparkSQL_hudi.ipynb Delta ãã¥ãŒããªã¢ã«ããŒãããã¯: s3://athena-examples-us-east-1/athenasparksqlblog/notebooks/SparkSQL_delta.ipynb ããŒãããã¯ãããŠã³ããŒããããã ããŒãããã¯ãã¡ã€ã«ã®ç®¡ç ïŒâ»æ³šïŒæ¬çš¿ç¿»èš³æç¹ã§ã¯ãªã³ã¯å
ã®ããã¥ã¡ã³ããæªç¿»èš³ã§ãããè±èªã§ã®æäŸã§ãã以äžåæ§ãïŒã® ããŒãããã¯ã®ã€ã³ããŒãæ¹æ³ ã®ã»ã¯ã·ã§ã³ã«åŸã£ãŠãAthena Spark ç°å¢ã«ã€ã³ããŒãããŠãã ããã Open Table Format ã«ããããã»ã¯ã·ã§ã³ãåç
§ããŠãã ãã Iceberg ããŒãã«åœ¢åŒã«èå³ãããå Žåã¯ã Apache Iceberg ããŒãã«ã®å©çš ã»ã¯ã·ã§ã³ãåç
§ããŠãã ããã Hudi ããŒãã«åœ¢åŒã«èå³ãããå Žåã¯ã Apache Hudi ããŒãã«ã®å©çš ã®ã»ã¯ã·ã§ã³ãåç
§ããŠãã ããã Delta Lake ããŒãã«åœ¢åŒã«èå³ãããå Žåã¯ã Linux Foundation Delta Lake ããŒãã«ã®å©çš ã®ã»ã¯ã·ã§ã³ãåç
§ããŠãã ããã Apache Iceberg ããŒãã«ã®å©çš Athena ã® Spark ããŒãããã¯ã䜿çšããéãPySpark ã®ã³ãŒãã䜿çšããããšãªãçŽæ¥ SQL ã¯ãšãªãå®è¡ã§ããŸãã ããã¯ãã»ã«ã®åäœã倿ŽããããŒãããã¯ã»ã«ã®ç¹å¥ãªãããã§ããã»ã«ããžãã¯ã䜿çšããããšã§å®çŸããŸãã SQL ã®å Žåã %%sql ããžãã¯ã远å ã§ããŸããããã«ãããã»ã«ã®å
容å
šäœã Athena äžã§å®è¡ããã SQL ã¹ããŒãã¡ã³ããšããŠè§£éãããŸãã ãã®ã»ã¯ã·ã§ã³ã§ã¯ãAthena ã® Apache Spark SQL ã䜿çšããŠãApache Iceberg ããŒãã«ãäœæãåæã管çããæ¹æ³ã瀺ããŸãã ããŒãããã¯ã»ãã·ã§ã³ã®èšå® Athena ã§ Apache Iceberg ã䜿çšããã«ã¯ãã»ãã·ã§ã³ã®äœæãç·šéäžã«ã Apache Spark ãããã㣠ã»ã¯ã·ã§ã³ãå±éãã Apache Iceberg ãªãã·ã§ã³ãéžæããŸãã 以äžã®ã¹ã¯ãªãŒã³ã·ã§ããã«ç€ºãããã«ãããããã£ãäºåã«èšå®ãããŸãã ã¹ãããã®è©³çްã¯ã ã»ãã·ã§ã³ã®è©³çްãç·šéãã ãŸã㯠ç¬èªã®ããŒãããã¯ãäœæãã ãã芧ãã ããã ãã®ã»ã¯ã·ã§ã³ã§äœ¿çšãããŠããã³ãŒãã¯ã SparkSQL_iceberg.ipynb ãã¡ã€ã«ã§åç
§ã§ããŸãã ããŒã¿ããŒã¹ãšIcebergããŒãã«ã®äœæ ãŸããAWS Glue ããŒã¿ã«ã¿ãã°ã«ããŒã¿ããŒã¹ãäœæããŸãã æ¬¡ã® SQL ã䜿çšãããšã icebergdb ãšããããŒã¿ããŒã¹ãäœæã§ããŸãã %%sql CREATE DATABASE icebergdb 次ã«ãããŒã¿ããŒã¹ icebergdb ã§ãããŒã¿ã®ããŒãå
ã«ãªã Amazon S3 ã®å Žæãæã noaa_iceberg ãšãã Iceberg ããŒãã«ãäœæããŸããæ¬¡ã®ã¹ããŒãã¡ã³ããå®è¡ãããã±ãŒã·ã§ã³ s3://<your-S3-bucket>/<prefix>/ ããèªèº«ã® S3 ãã±ãããšãã¬ãã£ãã¯ã¹ã«çœ®ãæããŠãã ãã: %%sql CREATE TABLE icebergdb.noaa_iceberg( station string, date string, latitude string, longitude string, elevation string, name string, temp string, temp_attributes string, dewp string, dewp_attributes string, slp string, slp_attributes string, stp string, stp_attributes string, visib string, visib_attributes string, wdsp string, wdsp_attributes string, mxspd string, gust string, max string, max_attributes string, min string, min_attributes string, prcp string, prcp_attributes string, sndp string, frshtt string) USING iceberg PARTITIONED BY (year string) LOCATION 's3://<your-S3-bucket>/<prefix>/noaaiceberg/' ããŒãã«ãžã®ããŒã¿ã®æ¿å
¥ noaa_iceberg Iceberg ããŒãã«ã«ããŒã¿ãå
¥åããããã«ãåææ¡ä»¶ãšããŠäœæããã Parquet ããŒãã« sparkblogdb.noaa_pq ããããŒã¿ãæ¿å
¥ããŸãã Spark ã® INSERT INTO ã¹ããŒãã¡ã³ãã䜿çšããŠãããè¡ãããšãã§ããŸãã %%sql INSERT INTO icebergdb.noaa_iceberg select * from sparkblogdb.noaa_pq ãããã¯ã CREATE TABLE AS SELECT ã« USING iceberg å¥ã䜿çšããããšã§ã Iceberg ããŒãã«ãäœæãããœãŒã¹ããŒãã«ããããŒã¿ãæ¿å
¥ããäžé£ã®ã¹ããããäžåºŠã«å®è¡ã§ããŸãã %%sql CREATE TABLE icebergdb.noaa_iceberg USING iceberg PARTITIONED BY (year) AS SELECT * FROM sparkblogdb.noaa_pq Iceberg ããŒãã«ã®ã¯ãšãª ããŒã¿ã Iceberg ããŒãã«ã«æ¿å
¥ãããã®ã§ãåæãéå§ã§ããŸãã 'SEATTLE TACOMA AIRPORT, WA US' ã®å Žæã«ã€ããŠã幎ããšã®æäœèšé²æ°æž©ãèŠã€ããããã«ãSpark SQL ãå®è¡ããŠã¿ãŸãããã %%sql select name, year, min(MIN) as minimum_temperature from icebergdb.noaa_iceberg where name = 'SEATTLE TACOMA AIRPORT, WA US' group by 1,2 次ã®ãããªåºåãåŸãããŸãã Iceberg ããŒãã«å
ã®ããŒã¿ã®æŽæ° ããŒãã«å
ã®ããŒã¿ãæŽæ°ããæ¹æ³ãèŠãŠãããŸãããã ã¹ããŒã·ã§ã³å 'SEATTLE TACOMA AIRPORT, WA US' ã 'Sea-Tac' ã«æŽæ°ããããšããŸãã Spark SQL ã䜿çšãããšãIceberg ããŒãã«ã«å¯Ÿã㊠UPDATE ã¹ããŒãã¡ã³ããå®è¡ã§ããŸãã %%sql UPDATE icebergdb.noaa_iceberg SET name = 'Sea-Tac' WHERE name = 'SEATTLE TACOMA AIRPORT, WA US' ãŸããåã®SELECTã¯ãšãªãå®è¡ããŠã 'Sea-Tac' ãã±ãŒã·ã§ã³ã®æäœèšé²æž©åºŠãèŠã€ããããšãã§ããŸãã %%sql select name, year, min(MIN) as minimum_temperature from icebergdb.noaa_iceberg where name = 'Sea-Tac' group by 1,2 次ã®ãããªåºåãåŸãããŸãã ããŒã¿ãã¡ã€ã«ã®å§çž® Icebergã®ãã㪠Open Table Format ã«ãããæŽæ°åŠçã§ã¯ããã¡ã€ã«ã¹ãã¬ãŒãžå
ã®æŽæ°å·®åãäœæãããããã§ã¹ããã¡ã€ã«ãéããŠè¡ã®ããŒãžã§ã³ããã©ããã³ã°ããããšã§ãã®æ©èœãå®çŸããŠããŸãã ãã¡ã€ã«æ°ãå€ããªããšãããã§ã¹ããã¡ã€ã«ã«æ ŒçŽãããã¡ã¿ããŒã¿ã®éãå€ããªããŸãããå°ããããŒã¿ã倧éã«ãããšäžå¿
èŠã«ã¡ã¿ããŒã¿ãå€ãããã¡ã§ãããã«ããã¯ãšãªå¹çã®äœäžãš Amazon S3 ã¢ã¯ã»ã¹ã³ã¹ãã®äžæãæããŸãã Athena (for Spark) ã§ Iceberg ã® rewrite_data_files ããã·ãŒãžã£ãå®è¡ãããšãããŒã¿ãã¡ã€ã«ãå§çž®ããã倿°ã®å°ããªå·®åãã¡ã€ã«ããèªã¿åãã«æé©åãããå°æ°ã® Parquet ãã¡ã€ã«ã«ãŸãšããããŸãã ãã¡ã€ã«ã®å§çž®ã«ããèªã¿åãæäœãé«éåãããŸãã ããŒãã«ã®å§çž®ãå®è¡ããã«ã¯ã次㮠Spark SQL ãå®è¡ããŸãã %%sql CALL spark_catalog.system.rewrite_data_files (table => 'icebergdb.noaa_iceberg', strategy=>'sort', sort_order => 'zorder(name)') rewrite_data_files ã«ã¯ ãœãŒãæŠç¥ãæå®ãããªãã·ã§ã³ããããããã«ããããŒã¿ã®åç·šæãšå§çž®ãé©åã«æå®ããããšãã§ããŸãã ããŒãã«ã¹ãããã·ã§ããã®ãªã¹ã衚瀺 Iceberg ããŒãã«äžã®åæžã蟌ã¿ãæŽæ°ãåé€ãã¢ãããµãŒããå§çž®æäœã¯ãã¹ãããã·ã§ããåé¢ãšã¿ã€ã ãã©ãã«ãå®çŸãããããå€ãããŒã¿ãšã¡ã¿ããŒã¿ãä¿æãã€ã€ãããŒãã«ã®æ°ããã¹ãããã·ã§ãããäœæããŸããIceberg ããŒãã«ã®ã¹ãããã·ã§ããäžèЧãåŸãã«ã¯ã次㮠Spark SQL ã¹ããŒãã¡ã³ããå®è¡ããŸãã %%sql SELECT * FROM spark_catalog.icebergdb.noaa_iceberg.snapshots å€ãã¹ãããã·ã§ããã®æéåã äžèŠã«ãªã£ãããŒã¿ãã¡ã€ã«ãåé€ããããŒãã«ã¡ã¿ããŒã¿ã®ãµã€ãºãå°ããä¿ã€ããã«ãæéãæå®ããã¹ãããã·ã§ããã®å®æçãªåé€ãæšå¥šãããŸããæéåãã§ã¯ãªãã¹ãããã·ã§ããã§ãŸã å¿
èŠãšãããŠãããã¡ã€ã«ã¯åé€ãããŸãããAthena for Spark ã§ã¯ã次㮠SQL ãå®è¡ããŠãç¹å®ã®ã¿ã€ã ã¹ã¿ã³ããããå€ã icebergdb.noaa_iceberg ããŒãã«ã®ã¹ãããã·ã§ããã®æéåããèšå®ã§ããŸãã %%sql CALL spark_catalog.system.expire_snapshots ('icebergdb.noaa_iceberg', TIMESTAMP '2023-11-30 00:00:00.000') timestamp ã®å€ã¯ yyyy-MM-dd HH:mm:ss.fff ã®åœ¢åŒã®æååã§æå®ãããŠããããšã«æ³šæããŠãã ããã åºåã¯åé€ãããããŒã¿ãã¡ã€ã«ãšã¡ã¿ããŒã¿ãã¡ã€ã«ã®æ°ãã«ãŠã³ããããã®ã«ãªããŸãã ããŒãã«ãšããŒã¿ããŒã¹ã®åé€ ãã®æŒç¿ã§äœ¿çšããIcebergããŒãã«ãšAmazon S3ã®é¢é£ããŒã¿ãã¯ãªãŒã³ã¢ããããã«ã¯ã次ã®Spark SQLãå®è¡ã§ããŸãã %%sql DROP TABLE icebergdb.noaa_iceberg PURGE 次㮠Spark SQL ãå®è¡ããŠãicebergdb ããŒã¿ããŒã¹ãåé€ããŸãã %%sql DROP DATABASE icebergdb Athena ã§ Spark ã䜿çšã㊠Iceberg ããŒãã«ã§å®è¡ã§ãããã¹ãŠã®æäœã®è©³çްã«ã€ããŠã¯ãIceberg ããã¥ã¡ã³ãã® Spark ã¯ãšãª ãš Spark ããã·ãŒãžã£ ãåç
§ããŠãã ããã Apache Hudi ããŒãã«ã®å©ç𠿬¡ã«ãAthena ã® Spark SQL ã䜿çšããŠãApache Hudi ããŒãã«ãäœæãåæã管çããæ¹æ³ã瀺ããŸãã ããŒãããã¯ã»ãã·ã§ã³ã®èšå® Athena ã§ Apache Hudi ã䜿çšããã«ã¯ãã»ãã·ã§ã³ã®äœæãŸãã¯ç·šéäžã«ã Apache Spark ãããã㣠ã»ã¯ã·ã§ã³ãå±éãã Apache Hudi ãªãã·ã§ã³ãéžæããŸãã ã¹ãããã®è©³çްã¯ã ã»ãã·ã§ã³ã®è©³çްãç·šéãã ãŸã㯠ç¬èªã®ããŒãããã¯ãäœæãã ãã芧ãã ããã ãã®ã»ã¯ã·ã§ã³ã§äœ¿çšãããŠããã³ãŒãã¯ã SparkSQL_hudi.ipynb ãã¡ã€ã«ã§å©çšã§ããŸãã以äžã®ã¹ãããã確èªããããã«ãå©çšãã ããã ããŒã¿ããŒã¹ãšHudiããŒãã«ã®äœæ ãŸããAWS Glue ããŒã¿ã«ã¿ãã°ã«æ ŒçŽããã hudidb ãšããããŒã¿ããŒã¹ãäœæããŸãããã®æ¬¡ã« Hudi ããŒãã«ãäœæããŸãã %%sql CREATE DATABASE hudidb Amazon S3 ã®ããŒã¿ãããŒãããå Žæãæã Hudi ããŒãã«ãäœæããŸãã ãã®ããŒãã«ã¯ã ã³ããŒãªã³ã©ã€ã åã§ããããšã«æ³šæããŠãã ããã ããã¯ããŒãã« DDL ã® type='cow' ã«ãã£ãŠå®çŸ©ãããŠããŸãã stationãšdate ãè€åãã©ã€ããªããŒãpreCombinedFieldãyearãšããŠå®çŸ©ããŸããã ãŸããããŒãã«ã¯ year ã§ããŒãã£ã·ã§ã³åãããŠããŸãã æ¬¡ã®ã¹ããŒãã¡ã³ããå®è¡ãããã±ãŒã·ã§ã³ s3://<your-S3-bucket>/<prefix>/ ããèªèº«ã® S3 ãã±ãããšãã¬ãã£ãã¯ã¹ã«çœ®ãæããŠãã ãã: %%sql CREATE TABLE hudidb.noaa_hudi( station string, date string, latitude string, longitude string, elevation string, name string, temp string, temp_attributes string, dewp string, dewp_attributes string, slp string, slp_attributes string, stp string, stp_attributes string, visib string, visib_attributes string, wdsp string, wdsp_attributes string, mxspd string, gust string, max string, max_attributes string, min string, min_attributes string, prcp string, prcp_attributes string, sndp string, frshtt string, year string) USING HUDI PARTITIONED BY (year) TBLPROPERTIES( primaryKey = 'station, date', preCombineField = 'year', type = 'cow' ) LOCATION ''s3://<your-S3-bucket>/<prefix>/noaahudi/' ããŒãã«ãžã®ããŒã¿ã®æ¿å
¥ Iceberg ãšåæ§ã«ãåã®ã¹ãããã§äœæãã sparkblogdb.noaa_pq ããŒãã«ããããŒã¿ãèªã¿åãããšã«ãã£ãŠããŒãã«ã«ããŒã¿ãå
¥åããããã«ã INSERT INTO ã¹ããŒãã¡ã³ãã䜿çšããŸãã %%sql INSERT INTO hudidb.noaa_hudi select * from sparkblogdb.noaa_pq Hudi ããŒãã«ã®ã¯ãšãª ããŒãã«ãäœæãããã®ã§ã 'SEATTLE TACOMA AIRPORT, WA US' ãã±ãŒã·ã§ã³ã«ãããæé«æ°æž©ãæ€çŽ¢ããã¯ãšãªãå®è¡ããŠã¿ãŸãããã %%sql select name, year, max(MAX) as maximum_temperature from hudidb.noaa_hudi where name = 'SEATTLE TACOMA AIRPORT, WA US' group by 1,2 Hudi ããŒãã«å
ã®ããŒã¿ã®æŽæ° ã¹ããŒã·ã§ã³å(name)ã 'SEATTLE TACOMA AIRPORT, WA US' ãã 'SeaâTac' ã«å€æŽããŸãããã Athena for Spark ã§ ã¢ããããŒã ã¹ããŒãã¡ã³ããå®è¡ããããšã§ã noaa_hudi ããŒãã«ã®ã¬ã³ãŒããæŽæ°ã§ããŸãã %%sql UPDATE hudidb.noaa_hudi SET name = 'Sea-Tac' WHERE name = 'SEATTLE TACOMA AIRPORT, WA US' ãSea-Tacããã±ãŒã·ã§ã³ã§èšé²ãããæé«æ°æž©ãæ€çŽ¢ããããã«ãåã®SELECTã®æ¡ä»¶å¥ãâSea-Tacâã«å€ããŠå®è¡ããŸãã %%sql select name, year, max(MAX) as maximum_temperature from hudidb.noaa_hudi where name = 'Sea-Tac' group by 1,2 ã¿ã€ã ãã©ãã«ã¯ãšãª SQL ã§ã¿ã€ã ãã©ãã«ã¯ãšãªã䜿çšããããšã§ãéå»ã®ããŒã¿ã¹ãããã·ã§ãããåæã§ããŸããäŸ: %%sql select name, year, max(MAX) as maximum_temperature from hudidb.noaa_hudi timestamp as of '2023-12-01 23:53:43.100' where name = 'SEATTLE TACOMA AIRPORT, WA US' group by 1,2 ãã®ã¯ãšãªã¯ãéå»ã®ç¹å®ã®æç¹ã§ã®ã·ã¢ãã«ç©ºæž¯ã®æ°æž©ããŒã¿ããã§ãã¯ããŸãã timestamp å¥ã䜿ãããšã§ãçŸåšã®ããŒã¿ã倿Žããããšãªãéå»ã«æ»ãããšãã§ããŸãã timestamp ã®å€ã¯ yyyy-MM-dd HH:mm:ss.fff ã®ãã©ãŒãããã®æååã§æå®ãããŠããããšã«æ³šæããŠãã ããã ã¯ã©ã¹ã¿ãªã³ã°ã«ããã¯ãšãªéåºŠã®æé©å Athena ã§ã®ã¯ãšãªããã©ãŒãã³ã¹ãæ¹åããããã«ãSpark SQL ã䜿çšã㊠Hudi ããŒãã«ã§ ã¯ã©ã¹ã¿ãªã³ã° ãå®è¡ã§ããŸãã %%sql CALL run_clustering(table => 'hudidb.noaa_hudi', order => 'name') ããŒãã«ã®ã³ã³ãã¯ã·ã§ã³(compaction) ã³ã³ãã¯ã·ã§ã³ã¯Hudiã«ç¹æã® Merge On Read (MOR) ããŒãã«ã§æ¡çšãããŠããããŒãã«ãµãŒãã¹ã§ãè¡ããŒã¹ã®ãã°ãã¡ã€ã«ããã®æŽæ°ã宿çã«å¯Ÿå¿ããåããŒã¹ã®ããŒã¹ãã¡ã€ã«ã«ããŒãžããããšã§ãããŒã¹ãã¡ã€ã«ã®æ°ããããŒãžã§ã³ãçæããŸããã³ã³ãã¯ã·ã§ã³ã¯ Copy On Write (COW) ããŒãã«ã«ã¯é©çšãããã MOR ããŒãã«ã«ã®ã¿é©çšãããŸãã Athena for Spark ã䜿çšã㊠MOR ããŒãã«ã®ã³ã³ãã¯ã·ã§ã³ãå®è¡ããã«ã¯ã次ã®ã¯ãšãªãå®è¡ã§ããŸãã %%sql CALL run_compaction(op => 'run', table => 'hudi_table_mor'); ããŒãã«ãšããŒã¿ããŒã¹ã®åé€ ä»¥äžã® Spark SQL ãå®è¡ããŠãäœæãã Hudi ããŒãã«ãšãAmazon S3 ã®å Žæã«é¢é£ä»ããããããŒã¿ãåé€ããŠãã ãã: %%sql DROP TABLE hudidb.noaa_hudi PURGE 次㮠Spark SQL ãå®è¡ããŠãããŒã¿ããŒã¹ hudidb ãåé€ããŸãã %%sql DROP DATABASE hudidb Athena ã§ Spark ã䜿çšã㊠Hudi ããŒãã«ã§å®è¡ã§ãããã¹ãŠã®æäœã«ã€ããŠã¯ãHudi ããã¥ã¡ã³ãã® SQL DDL ãš Procedures ãåç
§ããŠãã ããã Linux Foundation Delta Lake ããŒãã«ã®å©ç𠿬¡ã«ãAthena ã® Spark SQL ã䜿çšã㊠Delta Lake ããŒãã«ãäœæãåæã管çããæ¹æ³ã瀺ããŸãã ããŒãããã¯ã»ãã·ã§ã³ã®èšå® Athena ã§ Spark ã䜿çšã㊠Delta Lake ãå©çšããã«ã¯ãã»ãã·ã§ã³ã®äœæãŸãã¯ç·šéäžã«ã Apache Spark ãããã㣠ã»ã¯ã·ã§ã³ãå±éãã Linux Foundation Delta Lake ãéžæããŸãã ã¹ãããã®è©³çްã¯ã ã»ãã·ã§ã³ã®è©³çްãç·šéãã ãŸã㯠ç¬èªã®ããŒãããã¯ãäœæãã ãã芧ãã ããã ãã®ã»ã¯ã·ã§ã³ã§äœ¿çšãããŠããã³ãŒãã¯ã SparkSQL_delta.ipynb ãã¡ã€ã«ã§å©çšã§ããŸããããã«ãå©çšãã ããã ããŒã¿ããŒã¹ãšDelta LakeããŒãã«ã®äœæ ãã®ã»ã¯ã·ã§ã³ã§ã¯ãAWS Glue ããŒã¿ã«ã¿ãã°ã«ããŒã¿ããŒã¹ãäœæããŸãã æ¬¡ã® SQL ã䜿çšãããšã deltalakedb ãšããããŒã¿ããŒã¹ãäœæã§ããŸãã %%sql CREATE DATABASE deltalakedb 次ã«ãããŒã¿ããŒã¹ deltalakedb ã§ãããŒã¿ãããŒããã Amazon S3 ã®å Žæãæã noaa_delta ãšãã Delta Lake ããŒãã«ãäœæããŸããæ¬¡ã®ã¹ããŒãã¡ã³ããå®è¡ãããã±ãŒã·ã§ã³ s3://<your-S3-bucket>/<prefix>/ ããèªèº«ã® S3 ãã±ãããšãã¬ãã£ãã¯ã¹ã«çœ®ãæããŠãã ãã: %%sql CREATE TABLE deltalakedb.noaa_delta( station string, date string, latitude string, longitude string, elevation string, name string, temp string, temp_attributes string, dewp string, dewp_attributes string, slp string, slp_attributes string, stp string, stp_attributes string, visib string, visib_attributes string, wdsp string, wdsp_attributes string, mxspd string, gust string, max string, max_attributes string, min string, min_attributes string, prcp string, prcp_attributes string, sndp string, frshtt string) USING delta PARTITIONED BY (year string) LOCATION 's3://<your-S3-bucket>/<prefix>/noaadelta/' ããŒãã«ãžã®ããŒã¿ã®æ¿å
¥ åã®æçš¿ã§äœæãã sparkblogdb.noaa_pq ããŒãã«ããããŒã¿ãèªã¿åãããšã«ãããããŒãã«ã«å
¥åããããã« INSERT INTO ã¹ããŒãã¡ã³ãã䜿çšããŸãã %%sql INSERT INTO deltalakedb.noaa_delta select * from sparkblogdb.noaa_pq CREATE TABLE AS SELECT ã䜿çšããŠã1 ã€ã®ã¯ãšãªã§ Delta Lake ããŒãã«ãäœæãããœãŒã¹ããŒãã«ããããŒã¿ãæ¿å
¥ããããšãã§ããŸãã Delta Lake ããŒãã«ã®ã¯ãšãª Delta Lake ããŒãã«ã«ããŒã¿ãæ¿å
¥ãããã®ã§ãåæãéå§ããããšãã§ããŸãã 'SEATTLE TACOMA AIRPORT, WA US' ãã±ãŒã·ã§ã³ã®æäœèšé²æž©åºŠãèŠã€ããããã«ãSpark SQL ãå®è¡ããŸãããã %%sql select name, year, max(MAX) as minimum_temperature from deltalakedb.noaa_delta where name = 'SEATTLE TACOMA AIRPORT, WA US' group by 1,2 Deltaã¬ã€ã¯ããŒãã«å
ã®ããŒã¿ã®æŽæ° ã¹ããŒã·ã§ã³åã 'SEATTLE TACOMA AIRPORT, WA US' ãã 'SeaâTac' ã«å€æŽããŸãããã Athena ã® Spark äžã§ noaa_delta ããŒãã«ã®ã¬ã³ãŒããæŽæ°ãã UPDATE ã¹ããŒãã¡ã³ããå®è¡ã§ããŸãã %%sql UPDATE deltalakedb.noaa_delta SET name = 'Sea-Tac' WHERE name = 'SEATTLE TACOMA AIRPORT, WA US' åã®SELECTã¯ãšãªãå®è¡ããŠã 'Sea-Tac' ãã±ãŒã·ã§ã³ã®æäœèšé²æž©åºŠãæ€çŽ¢ã§ããŸããçµæã¯ä»¥åãšåãã¯ãã§ãã %%sql select name, year, max(MAX) as minimum_temperature from deltalakedb.noaa_delta where name = 'Sea-Tac' group by 1,2 ããŒã¿ãã¡ã€ã«ã®å§çž® (optimize) Athena for Spark ã§ã¯ãDelta Lake ããŒãã«ã«å¯Ÿã㊠OPTIMIZE ãå®è¡ã§ããŸããããã«ãããè€æ°ã®å°ãããã¡ã€ã«ã倧ããªãã¡ã€ã«ã«å§çž®ããããããå°ãããã¡ã€ã«ãããããããããšã«ããã¯ãšãªãžã®è² æ
ãæžããããšãã§ããŸããå§çž®ãå®è¡ããã«ã¯ã次ã®ã¯ãšãªãå®è¡ããŸãã %%sql OPTIMIZE deltalakedb.noaa_delta Delta Lake ã®ããã¥ã¡ã³ãã® æé©å ãåç
§ããŠãOPTIMIZE ã®å®è¡äžã«äœ¿çšã§ããããŸããŸãªãªãã·ã§ã³ã確èªããŠãã ããã Delta Lake ããŒãã«ã§åç
§ãããªããªã£ããã¡ã€ã«ã®åé€ Athena ã® Spark ã䜿çšã㊠Delta Lake ããŒãã«äžã§ VACUUM ã³ãã³ããå®è¡ããããšã§ããã®ããŒãã«ããåç
§ãããªããªã£ã Amazon S3 ã«ä¿åããããã¡ã€ã«ã§ãä¿ææéãè¶
ãããã®ãåé€ã§ããŸãã %%sql VACUUM deltalakedb.noaa_delta Delta Lake ã®ããã¥ã¡ã³ãã® Delta ããŒãã«ã§åç
§ãããªããªã£ããã¡ã€ã«ã®åé€ ãåç
§ããŠãVACUUM ã§å©çšã§ãããªãã·ã§ã³ã確èªããŠãã ããã ããŒãã«ãšããŒã¿ããŒã¹ã®åé€ æ¬¡ã® Spark SQL ãå®è¡ããŠãäœæãã Delta Lake ããŒãã«ãåé€ããŸãã %%sql DROP TABLE deltalakedb.noaa_delta 次㮠Spark SQL ãå®è¡ããŠãããŒã¿ããŒã¹ deltalakedb ãåé€ããŸãã %%sql DROP DATABASE deltalakedb Delta Lake ããŒãã«ãšããŒã¿ããŒã¹ã§ DROP TABLE DDL ãå®è¡ãããšããããã®ãªããžã§ã¯ãã®ã¡ã¿ããŒã¿ãåé€ãããŸãããAmazon S3 ã®ããŒã¿ãã¡ã€ã«ã¯èªåçã«ã¯åé€ãããŸããã ããŒãããã¯ã®ã»ã«ã§æ¬¡ã® Python ã³ãŒããå®è¡ããããšã§ãS3 ãã±ããããããŒã¿ãåé€ã§ããŸãã import boto3 s3 = boto3.resource('s3') bucket = s3.Bucket('') bucket.objects.filter(Prefix="/noaadelta/").delete() Athena ã® Spark ã䜿çšã㊠Delta Lake ããŒãã«ã§å®è¡ã§ãã SQL ã¹ããŒãã¡ã³ãã®è©³çްã«ã€ããŠã¯ãDelta Lake ããã¥ã¡ã³ãã® ã¯ã€ãã¯ã¹ã¿ãŒã ãåç
§ããŠãã ããã ãŸãšã ãã®æçš¿ã§ã¯ãAthena ããŒãããã¯ã§ Spark SQL ã䜿çšããŠããŒã¿ããŒã¹ãšããŒãã«ãäœæããããŒã¿ãæ¿å
¥ããã³ã¯ãšãªãå®è¡ããHudiãDelta LakeãIceberg ããŒãã«ã§ã®æŽæ°ãå§çž®ãã¿ã€ã ãã©ãã«ãªã©ã®äžè¬çãªæäœãå®è¡ããæ¹æ³ã瀺ããŸããã Open Table Format (OTF) ã¯ã ACID ãã©ã³ã¶ã¯ã·ã§ã³ãã¢ãããµãŒããåé€ãšãã£ãæäœãããŒã¿ã¬ã€ã¯äžã§å¯èœã«ãããªããžã§ã¯ãã¹ãã¬ãŒãžäžã§ã®ããé«åºŠãªæäœãæäŸããŸãã Athena for Spark ã§ã¯ãå¥éã³ãã¯ã¿ãã€ã³ã¹ããŒã«ããå¿
èŠããªãã®ã§ãAmazon S3 äžã«ä¿¡é Œã§ããããŒã¿ã¬ã€ã¯ãæ§ç¯ãéã®ãããã®äžè¬çãªæºåã管çã®ãªãŒããŒããããåæžããããšãã§ããŸãã ããŒã¿ã¬ã€ã¯äžã®åŠçã«ããã Open Table Format ã®éžæã®è©³çްã«ã€ããŠã¯ã AWS ã§ã®ãã©ã³ã¶ã¯ã·ã§ã³ããŒã¿ã¬ã€ã¯ã®ããã®ãªãŒãã³ããŒãã«ãã©ãŒãããã®éžæ ãåç
§ããŠãã ããã èè
ã«ã€ã㊠Pathik Shah ã¯ãAmazon Athena ã®ã·ãã¢ã¢ããªãã£ã¯ã¹ã¢ãŒããã¯ãã§ãã2015幎ã«AWSã«å å
¥ããŠä»¥æ¥ãããã°ããŒã¿ã¢ããªãã£ã¯ã¹ã®åéã«æ³šåããAWSã®ã¢ããªãã£ã¯ã¹ãµãŒãã¹ã䜿çšããŠã¹ã±ãŒã©ãã«ã§å
ç¢ãªãœãªã¥ãŒã·ã§ã³ã®æ§ç¯ãæ¯æŽããŠããŸãã Raj Devnath ã¯ãAmazon Athena ã®ãããã¯ããããŒãžã£ãŒã§ããã客æ§ã«æããã補åã®æ§ç¯ãšãã客æ§ã®ããŒã¿ãã䟡å€ãåŒãåºãããšã«æ
ç±ã泚ãã§ããŸããéèãå°å£²ãã¹ããŒããã«ãå®¶åºãªãŒãã¡ãŒã·ã§ã³ãããŒã¿éä¿¡ã·ã¹ãã ãªã©ãè€æ°ã®ãšã³ãããŒã±ããåãã®ãœãªã¥ãŒã·ã§ã³ã®æäŸçµéšããããŸãã 翻蚳ïŒãœãªã¥ãŒã·ã§ã³ã¢ãŒããã¯ã äžäœç² æ ( twitter â @simosako) åæïŒ Use Amazon Athena with Spark SQL for your open-source transactional table formats