æ¬èšäºã¯ã2025/5/9 ã«å
¬éããã Accelerate lightweight analytics using PyIceberg with AWS Lambda and an AWS Glue Iceberg REST endpoint ã翻蚳ãããã®ã§ãã翻蚳㯠Solutions Architect ã®æ·±èŠãæ
åœããŸããã ããŒã¿ã€ã³ãµã€ãã«åºã¥ã決å®ãè¡ãçŸä»£ã®çµç¹ã«ãšã£ãŠã广çãªããŒã¿ç®¡çã¯ãé«åºŠãªåæãšå¹ççãªæ©æ¢°åŠç¿ã®å©çšãå®çŸããããã®éèŠãªèŠçŽ ã§ããããŒã¿å©çšã®ãŠãŒã¹ã±ãŒã¹ãããè€éã«ãªãã«ã€ããããŒã¿ãšã³ãžãã¢ãªã³ã°ããŒã ã«ã¯ãè€æ°ã®ããŒã¿ãœãŒã¹ãã¢ããªã±ãŒã·ã§ã³å
šäœã§ã®ããŒãžã§ã³ç®¡çãå¢å ããããŒã¿éãã¹ããŒã倿Žã«å¯ŸåŠããããã®é«åºŠãªããŒã«ãå¿
èŠã«ãªããŸãã Apache Iceberg ã¯ãããŒã¿ã¬ã€ã¯ã§äººæ°ã®éžæè¢ãšãªã£ãŠããŸããACID (ååæ§ãäžè²«æ§ãç¬ç«æ§ãæ°žç¶æ§) ãã©ã³ã¶ã¯ã·ã§ã³ãã¹ããŒãé²åãã¿ã€ã ãã©ãã«æ©èœãæäŸããŸããIceberg ããŒãã«ã¯ãApache Spark ã Trino ãªã©ã®æ§ã
ãªåæ£ããŒã¿åŠçãã¬ãŒã ã¯ãŒã¯ããã¢ã¯ã»ã¹ã§ããããã倿§ãªããŒã¿åŠçã®ããŒãºã«å¯ŸããŠæè»ãªãœãªã¥ãŒã·ã§ã³ãšãªããŸãããã®ãã㪠Iceberg ãæ±ãããã®ããŒã«ã®äžã§ã PyIceberg ã¯åæ£ã³ã³ãã¥ãŒãã£ã³ã°ãªãœãŒã¹ãå¿
èŠãšããã«ãPython ã¹ã¯ãªããäžã§ããŒãã«ã®ã¢ã¯ã»ã¹ãšç®¡çãå¯èœã«ããŸãã ãã®æçš¿ã§ã¯ã AWS Glue Data Catalog ãš AWS Lambda ãšçµ±åããã PyIceberg ããçŽæç㪠Python ã€ã³ã¿ãŒãã§ãŒã¹ãéã㊠Iceberg ã®åŒ·åãªæ©èœã掻çšããããã®è»œéãªã¢ãããŒããæäŸããæ¹æ³ã瀺ããŸãããã®çµ±åã«ãããããŒã ã¯ã»ãšãã©ã»ããã¢ãããã€ã³ãã©ã¹ãã©ã¯ãã£ã®äŸåé¢ä¿ã®èšå®ãè¡ãããšã Iceberg ããŒãã«ã®æäœãå©çšãéå§ã§ããããšã説æããŸãã PyIceberg ã®äž»èŠæ©èœãšå©ç¹ PyIceberg ã®äž»ãªå©ç¹ã® 1 ã€ã¯ã軜éã§ããããšã§ãã忣ã³ã³ãã¥ãŒãã£ã³ã°ãã¬ãŒã ã¯ãŒã¯ãå¿
èŠãšãããããŒã 㯠Python ã¢ããªã±ãŒã·ã§ã³ããçŽæ¥ããŒãã«æäœãå®è¡ã§ãããããåŠç¿æ²ç·ãå°ãããå°èŠæš¡ããäžèŠæš¡ã®ããŒã¿æ¢çŽ¢ãšåæã«é©ããŠããŸããããã«ãPyIceberg 㯠Pandas ã Polars ãªã©ã® Python ããŒã¿åæã©ã€ãã©ãªãšçµ±åãããŠãããããããŒã¿ãŠãŒã¶ãŒã¯æ¢åã®ã¹ãã«ãšã¯ãŒã¯ãããŒã掻çšã§ããŸãã PyIceberg ã Data Catalog ãš Amazon Simple Storage Service (Amazon S3) ã§äœ¿çšãããšãããŒã¿ããŒã ã¯ããŒãã«ãå®å
šã«ãµãŒããŒã¬ã¹ãªç°å¢ã§å©çšããã³ç®¡çã§ããŸããã€ãŸããããŒã¿ããŒã ã¯ã€ã³ãã©ã¹ãã©ã¯ãã£ã®ç®¡çã§ã¯ãªããåæãšæŽå¯ã«éäžããããšãã§ããŸãã ããã«ãPyIceberg ãéããŠç®¡çããã Iceberg ããŒãã«ã¯ãAWS ã®ããŒã¿åæãµãŒãã¹ãšäºææ§ããããŸããPyIceberg ã¯åâŒããŒãã§åäœããããã⌀éã®ããŒã¿ãæ±ãå Žåã¯æ§èœã«å¶éããããŸããã Amazon Athena ã AWS Glue ãªã©ã®ãµãŒãã¹ã䜿ãã°ãåãããŒãã«ãããå€§èŠæš¡ã«å¹ççã«åŠçã§ããŸããããã«ãããããŒã 㯠PyIceberg ã䜿ã£ãŠè¿
éãªéçºãšãã¹ããè¡ãããã®åŸãããŒã¿ç®¡çã¢ãããŒãã®äžè²«æ§ãç¶æããªãããããå€§èŠæš¡ãªåŠçãšã³ãžã³ã䜿ã£ãæ¬çªã¯ãŒã¯ããŒãã«ã·ãŒã ã¬ã¹ã«ç§»è¡ã§ããŸãã 代衚çãªãŠãŒã¹ã±ãŒã¹ 次ã®ãããªã·ããªãªã§ã¯ãPyIceberg ãç¹ã«åœ¹ç«ã¡ãŸã: ããŒã¿ãµã€ãšã³ã¹ã®å®éšãšç¹åŸŽéãšã³ãžãã¢ãªã³ã° â ããŒã¿ãµã€ãšã³ã¹ã§ã¯ãä¿¡é Œã§ããå¹ççãªåæãšã¢ãã«ãç¶æããããã«ãå®éšã®åçŸæ§ãéèŠã§ããããããçµç¹ã®ããŒã¿ãç¶ç¶çã«æŽæ°ããããããéèŠãªããžãã¹ã€ãã³ããã¢ãã«åŠç¿ãäžè²«ããåç
§ã®ããã®ããŒã¿ã¹ãããã·ã§ããã管çããããšãé£ãããªããŸããããŒã¿ãµã€ãšã³ãã£ã¹ãã¯ãã¿ã€ã ãã©ãã«æ©èœã䜿ã£ãŠããŒã¿ã®éå»ã®ã¹ãããã·ã§ãããç
§äŒãã ã¿ã°ä»ãæ©èœ ã䜿ã£ãŠéèŠãªããŒãžã§ã³ãèšé²ã§ããŸããPyIceberg ã䜿ãã°ãPandas ãªã©ã®éŠŽæã¿ã®ããããŒã«ã䜿ã£ãŠ Python ç°å¢ã§ãããã®å©ç¹ãåŸãããŸããIceberg ã® ACID ç¹æ§ã®ãããã§ãããŒãã«ãç©æ¥µçã«æŽæ°ãããŠããå Žåã§ãæŽåæ§ãæ
ä¿ããããŒã¿ã¢ã¯ã»ã¹ãå¯èœã«ãªããŸãã AWS Lambda ã«ãããµãŒããŒã¬ã¹ããŒã¿åŠç â çµç¹ã¯å€ãã®å Žåãè€éãªã€ã³ãã©ã¹ãã©ã¯ãã£ã管çããã«ãããŒã¿ãå¹ççã«åŠçããåæããŒãã«ãç¶æããå¿
èŠããããŸããPyIceberg ãš Lambda ã䜿ãã°ãããŒã ã¯ãµãŒããŒã¬ã¹ãª Lamnba 颿°ã䜿ã£ãŠã€ãã³ãé§åã®ããŒã¿åŠçãã¹ã±ãžã¥ãŒã«ãããããŒãã«æŽæ°ãæ§ç¯ã§ããŸããPyIceberg ã®è»œéãªæ§è³ªã¯ããµãŒããŒã¬ã¹ç°å¢ã«é©ããŠãããããŒã¿æ€èšŒã倿ãåã蟌ã¿ãªã©ã®ã·ã³ãã«ãªããŒã¿åŠçã¿ã¹ã¯ãå¯èœã«ããŸãããããã®ããŒãã«ã¯ãããŸããŸãª AWS ãµãŒãã¹ãéããŠæŽæ°ãšåæã®äž¡æ¹ã«ã¢ã¯ã»ã¹ã§ãããããããŒã ã¯ãµãŒããŒãã¯ã©ã¹ã¿ãŒã管çããã«å¹ççãªããŒã¿ãã€ãã©ã€ã³ãæ§ç¯ã§ããŸãã PyIceberg ã䜿çšããã€ãã³ãé§åã®ããŒã¿åã蟌ã¿ãšåæ ãã®ã»ã¯ã·ã§ã³ã§ã¯ã NYC yellow taxi trip data ã䜿çšããŠãPyIceberg ã«ããããŒã¿åŠçãšåæã®å®è·µçãªäŸãæ¢ããŸãããªã¢ã«ã¿ã€ã ã®ã¿ã¯ã·ãŒèµ°è¡èšé²ã®åŠçãã·ãã¥ã¬ãŒãããããã«ãLambda ã䜿çšããŠãµã³ãã«ããŒã¿ã Iceberg ããŒãã«ã«æ¿å
¥ããŸãããã®äŸã§ã¯ãå¹ççãªããŒã¿åã蟌ã¿ãšæè»ãªåææ©èœãçµã¿åãããããšã§ãã¯ãŒã¯ãããŒãããåççãªãã®ã«ããæ¹æ³ã瀺ããŸãã ããŒã ãæ¬¡ã®ãããªè€æ°ã®èŠä»¶å¯Ÿå¿ããå¿
èŠãããå Žé¢ãæ³åããŠãã ãã: ããŒã¿åŠçãœãªã¥ãŒã·ã§ã³ã¯ãäžèŠæš¡ã®ããŒã¿ã»ãããåŠçããã±ãŒã¹ã§åæ£ã³ã³ãã¥ãŒãã£ã³ã°ã¯ã©ã¹ã¿ãŒã®ç®¡çã®è€éããé¿ãããããã³ã¹ãããã©ãŒãã³ã¹ãè¯ããã¡ã³ããã³ã¹æ§ãé«ãå¿
èŠããããŸãã ã¢ããªã¹ãã¯ãPython ã®ããŒã«ã䜿ã£ãŠæè»ãªã¯ãšãªãšæ¢çŽ¢ãè¡ããããã«ããå¿
èŠããããŸããäŸãã°ãéå»ã®ã¹ãããã·ã§ãããšçŸåšã®ããŒã¿ãæ¯èŒããŠãæéã®çµéã«äŒŽããã¬ã³ããåæããå¿
èŠããããããããŸããã ãœãªã¥ãŒã·ã§ã³ã¯ãå°æ¥çã«ããæ¡åŒµæ§ã®é«ããã®ã«ããèœåãæã€å¿
èŠããããŸãã ãããã®èŠä»¶ã«å¯ŸåŠãããããPyIceberg ã§åäœãã Lambda ã«ããããŒã¿åŠçãš Jupyter ããŒãããã¯ã«ããåæãçµã¿åããããœãªã¥ãŒã·ã§ã³ãå®è£
ããŸãããã®ã¢ãããŒãã«ãããããŒã¿ã®æŽåæ§ãç¶æããªããæè»ãªåæã¯ãŒã¯ãããŒãå¯èœã«ããã軜éã§ãããªããå
ç¢ãªã¢ãŒããã¯ãã£ãå®çŸãããŸãããŠã©ãŒã¯ã¹ã«ãŒã®æåŸã§ã¯ãAthena ã䜿çšããŠãã®ããŒã¿ãç
§äŒããè€æ°ã® Iceberg 察å¿ããŒã«ãšã®äºææ§ã瀺ããšãšãã«ããã®ã¢ãŒããã¯ãã£ã®ã¹ã±ãŒã«æ§ã瀺ããŸãã 倧ãŸããªæé ã¯ä»¥äžã®éãã§ã: Lambda ã䜿çšããŠã AWS Glue Iceberg REST ãšã³ããã€ã³ã çµç±ã§ PyIceberg ãå©çšãããµã³ãã«ã® NYC yellow taxi trip data ã Amazon S3 äžã® Iceberg ããŒãã«ã«æžã蟌ã¿ãŸããå®éã®ã·ããªãªã§ã¯ããã® Lambda 颿°ã¯ Amazon Simple Queue Service (Amazon SQS) ãªã©ã®ãã¥ãŒã€ã³ã°ã³ã³ããŒãã³ãããã®ã€ãã³ãã«ãã£ãŠããªã¬ãŒãããŸãã詳现ã¯ã Lambda ãš Amazon SQS ã®äœµçš ãåç
§ããŠãã ããã Jupyter ããŒãããã¯ã§ AWS Glue Iceberg REST ãšã³ããã€ã³ããçµç±ã㊠PyIceberg ã䜿çšããããŒãã«ããŒã¿ãåæããŸãã Athena ã䜿çšããŠããŒã¿ãã¯ãšãªããIceberg ã®æè»æ§ãå®èšŒããŸãã æ¬¡ã®å³ã¯ã¢ãŒããã¯ãã£ã瀺ããŠããŸãã ãã®ã¢ãŒããã¯ãã£ãå®è£
ããéã«éèŠã«ãªãç¹ã¯ãLambda 颿°ãã€ãã³ãã«ãã£ãŠããªã¬ãŒããããšãè€æ°ã®åæå®è¡ãçºçããå¯èœæ§ãããããšã§ãããã®åæå®è¡ã«ãããIceberg ããŒãã«ãžã®æžãèŸŒã¿æã«ãã©ã³ã¶ã¯ã·ã§ã³ç«¶åãçºçããå¯èœæ§ããããŸãããããåŠçããã«ã¯ãé©åãªå詊è¡ã¡ã«ããºã ãå®è£
ãããã©ã³ã¶ã¯ã·ã§ã³åé¢ã¬ãã«ãæ
éã«ç®¡çããå¿
èŠããããŸããAmazon SQS ãã€ãã³ããœãŒã¹ãšããŠäœ¿çšããå Žåã¯ã SQS ã€ãã³ããœãŒã¹ã®æå€§åæå®è¡èšå® ã䜿ã£ãŠåæå®è¡ãå¶åŸ¡ã§ããŸãã åææ¡ä»¶ ãã®ãŠãŒã¹ã±ãŒã¹ã«ã¯ã次ã®åææ¡ä»¶ãå¿
èŠã§ãã LambdaãAWS GlueãAmazon S3ãAthenaãããã³ AWS CloudFormation ãžã®ã¢ã¯ã»ã¹ãæäŸããã¢ã¯ãã£ã㪠AWS ã¢ã«ãŠã³ãã CloudFormation ã¹ã¿ãã¯ãäœæããã³ãããã€ããæš©éã詳现ã«ã€ããŠã¯ã èªå·±ç®¡çæš©éã䜿çšãã CloudFormation StackSets ã®äœæ ãåç
§ããŠãã ããã AWS CloudShell ãšãã®æ©èœã«å¯Ÿããå®å
šãªã¢ã¯ã»ã¹æš©ãæã€ãŠãŒã¶ãŒã詳现ã«ã€ããŠã¯ã AWS CloudShell ã®æŠèŠ ãåç
§ããŠãã ããã AWS CloudFormation ã«ãããªãœãŒã¹ã®ã»ããã¢ãã æ¬¡ã® CloudFormation ãã³ãã¬ãŒãã䜿çšããŠã以äžã®ãªãœãŒã¹ãã»ããã¢ããã§ããŸã: Iceberg ããŒãã«ã®ã¡ã¿ããŒã¿ãšããŒã¿ãã¡ã€ã«ãæ ŒçŽãã S3 ãã±ãã Amazon Elastic Container Registry (Amazon ECR) ã®ãªããžã㪠(Lambda 颿°ã®ã³ã³ããã€ã¡ãŒãžãæ ŒçŽ) ããŒãã«ãæ ŒçŽããããŒã¿ã«ã¿ãã°ããŒã¿ããŒã¹ Amazon SageMaker AI ã® ããŒãããã¯ã€ã³ã¹ã¿ã³ã¹ (Jupyter ããŒãããã¯ç°å¢çš) Lambda 颿°ãš SageMaker AI ããŒãããã¯ã€ã³ã¹ã¿ã³ã¹çšã® AWS Identity and Access Management (IAM) ããŒã« æ¬¡ã®æé ã«åŸã£ãŠãªãœãŒã¹ããããã€ããŠãã ããã Launch stack ãéžæããŸãã ã¹ã¿ãã¯ã® ãã©ã¡ãŒã¿ ã§ã¯ãããŒã¿ããŒã¹åãšããŠã pyiceberg_lambda_blog_database ãããã©ã«ãã§èšå®ãããŠããŸããããã©ã«ãå€ã倿Žããããšãã§ããŸããããŒã¿ããŒã¹åã倿Žããå Žåã¯ã以éã®ãã¹ãŠã®ã¹ãããã§ pyiceberg_lambda_blog_database ãéžæããååã«çœ®ãæããããšãå¿ããªãã§ãã ãããæ¬¡ã«ã 次㞠ãéžæããŸãã æ¬¡ãž ãéžæããŸãã AWS CloudFormation ã«ãã£ãŠ IAM ãªãœãŒã¹ãã«ã¹ã¿ã åã§äœæãããå Žåãããããšãæ¿èªããŸã ãéžæããŸãã éä¿¡ ãéžæããŸãã Lambda 颿°ã®æ§ç¯ãšå®è¡ PyIceberg ã䜿ã£ãŠçä¿¡ã¬ã³ãŒããåŠçãã Lambda 颿°ãæ§ç¯ããŸãããããã®é¢æ°ã¯ãData Catalog ã® pyiceberg_lambda_blog_database ããŒã¿ããŒã¹å
ã« nyc_yellow_table ãšãã Iceberg ããŒãã«ãååšããªãå Žåã«æ°èŠã§ããŒãã«ãäœæããŸããæ¬¡ã«ããµã³ãã«ã® NYC yellow taxi trip data ãçæããŠãã¬ã³ãŒããã·ãã¥ã¬ãŒããã nyc_yellow_table ã«æ¿å
¥ããŸãã ãã®äŸã§ã¯ããã®é¢æ°ãæåã§åŒã³åºããŠããŸãããå®éã®ã·ããªãªã§ã¯ããã® Lambda 颿°ã¯ Amazon SQS ããã®ã¡ãã»ãŒãžãªã©ã®å®éã®ã€ãã³ãã«ãã£ãŠããªã¬ãŒãããŸããå®éã®ãŠãŒã¹ã±ãŒã¹ãå®è£
ããéã¯ã颿°ã³ãŒããã€ãã³ãããŒã¿ãåãåããèŠä»¶ã«åºã¥ããŠåŠçããããã«å€æŽããå¿
èŠããããŸãã ã³ã³ããã€ã¡ãŒãžããããã€ããã±ãŒãžãšããŠäœ¿çšããŠããã®é¢æ°ããããã€ããŸããããã§ã¯ã³ã³ããã€ã¡ãŒãžãã Lambda 颿°ãäœæããããã«ãCloudShell ã§ã€ã¡ãŒãžããã«ãããECR ãªããžããªã«ããã·ã¥ããŸãã以äžã®æé ãå®è¡ããŠãã ããã AWS ãããžã¡ã³ãã³ã³ãœãŒã« ã«ãµã€ã³ã€ã³ãã CloudShell ãèµ·å ããŸãã äœæ¥ãã£ã¬ã¯ããªãäœæããŸãã mkdir pyiceberg_blog cd pyiceberg_blog Lambda ã¹ã¯ãªãã lambda_function.py ãããŠã³ããŒãããŸãã aws s3 cp s3://aws-blogs-artifacts-public/artifacts/BDB-5013/lambda_function.py . ãã®ã¹ã¯ãªããã¯ä»¥äžã®ã¿ã¹ã¯ãå®è¡ããŸã: Data Catalog ã« NYC yellow taxi trip data ã®ã¹ããŒããæã€ Iceberg ããŒãã«ãäœæããŸã ã©ã³ãã 㪠NYC yellow taxi trip data ã«ããšã¥ãããŒã¿ã»ãããçæããŸã ãã®ããŒã¿ãããŒãã«ã«æ¿å
¥ããŸã ãã® Lambda 颿°ã®éèŠãªéšåãæãäžããŠã¿ãŸããã: Iceberg ã«ã¿ãã°ã®æ§æ â æ¬¡ã®ã³ãŒãã¯ãAWS Glue Iceberg REST ãšã³ããã€ã³ãã«æ¥ç¶ãã Iceberg ã«ã¿ãã°ãå®çŸ©ããŠããŸã: # Configure the catalog catalog_properties = { "type": "rest", "uri": f"https://glue.{region}.amazonaws.com/iceberg", "s3.region": region, "rest.sigv4-enabled": "true", "rest.signing-name": "glue", "rest.signing-region": region } catalog = load_catalog(**catalog_properties) ããŒãã«ã¹ããŒãå®çŸ© â æ¬¡ã®ã³ãŒãã¯ãNYC taxi ããŒã¿ã»ããã® Iceberg ããŒãã«ã¹ããŒããå®çŸ©ããŠããŸãããã®ããŒãã«ã«ã¯ä»¥äžãå«ãŸããŠããŸãã Schema ã§å®çŸ©ãããã¹ããŒãã®ã«ã©ã vendorid ãš tpep_pickup_datetime ã«ãã ããŒãã£ã·ã§ãã³ã° (PartitionSpec ã䜿çš) tpep_pickup_datetime ã«é©çšããã day 倿 tpep_pickup_datetime ãš tpep_dropoff_datetime ã«ãã ãœãŒããªãŒã㌠ã¿ã€ã ã¹ã¿ã³ãåã« day 倿ãé©çšããéãIceberg ã¯éå±€çã«æ¥ä»ããŒã¹ã®ããŒãã£ã·ã§ã³åå²ãèªåçã«åŠçããŸããã€ãŸããåäžã® day 倿ã§ã幎ãæãæ¥ã®ã¬ãã«ã§ããŒãã£ã·ã§ã³ãã«ãŒãã³ã°ãå¯èœã«ãããããåã¬ãã«ã«æç€ºçãªå€æãå¿
èŠãšããŸãããIceberg ã®ããŒãã£ã·ã§ã³åå²ã®è©³çްã«ã€ããŠã¯ã ããŒãã£ã·ã§ã³åå² ãåç
§ããŠãã ããã # Table Definition schema = Schema( NestedField(field_id=1, name="vendorid", field_type=LongType(), required=False), NestedField(field_id=2, name="tpep_pickup_datetime", field_type=TimestampType(), required=False), NestedField(field_id=3, name="tpep_dropoff_datetime", field_type=TimestampType(), required=False), NestedField(field_id=4, name="passenger_count", field_type=LongType(), required=False), NestedField(field_id=5, name="trip_distance", field_type=DoubleType(), required=False), NestedField(field_id=6, name="ratecodeid", field_type=LongType(), required=False), NestedField(field_id=7, name="store_and_fwd_flag", field_type=StringType(), required=False), NestedField(field_id=8, name="pulocationid", field_type=LongType(), required=False), NestedField(field_id=9, name="dolocationid", field_type=LongType(), required=False), NestedField(field_id=10, name="payment_type", field_type=LongType(), required=False), NestedField(field_id=11, name="fare_amount", field_type=DoubleType(), required=False), NestedField(field_id=12, name="extra", field_type=DoubleType(), required=False), NestedField(field_id=13, name="mta_tax", field_type=DoubleType(), required=False), NestedField(field_id=14, name="tip_amount", field_type=DoubleType(), required=False), NestedField(field_id=15, name="tolls_amount", field_type=DoubleType(), required=False), NestedField(field_id=16, name="improvement_surcharge", field_type=DoubleType(), required=False), NestedField(field_id=17, name="total_amount", field_type=DoubleType(), required=False), NestedField(field_id=18, name="congestion_surcharge", field_type=DoubleType(), required=False), NestedField(field_id=19, name="airport_fee", field_type=DoubleType(), required=False), ) # Define partition spec partition_spec = PartitionSpec( PartitionField(source_id=1, field_id=1001, transform=IdentityTransform(), name="vendorid_idenitty"), PartitionField(source_id=2, field_id=1002, transform=DayTransform(), name="tpep_pickup_day"), ) # Define sort order sort_order = SortOrder( SortField(source_id=2, transform=DayTransform()), SortField(source_id=3, transform=DayTransform()) ) database_name = os.environ.get('GLUE_DATABASE_NAME') table_name = os.environ.get('ICEBERG_TABLE_NAME') identifier = f"{database_name}.{table_name}" # Create the table if it doesn't exist location = f"s3://pyiceberg-lambda-blog-{account_id}-{region}/{database_name}/{table_name}" if not catalog.table_exists(identifier): table = catalog.create_table( identifier=identifier, schema=schema, location=location, partition_spec=partition_spec, sort_order=sort_order ) else: table = catalog.load_table(identifier=identifier) ããŒã¿ã®çæãšæ¿å
¥ â æ¬¡ã®ã³ãŒãã¯ã©ã³ãã ãªããŒã¿ãçæããããŒãã«ã«æ¿å
¥ããŸãããã®äŸã§ã¯ãããžãã¹ã€ãã³ãããã©ã³ã¶ã¯ã·ã§ã³ã远跡ããããã«ãæ°ããã¬ã³ãŒããç¶ç¶çã«è¿œå ããã append-only ã®ãã¿ãŒã³ãæ³å®ããŠããŸãã # Generate random data records = generate_random_data() # Convert to Arrow Table df = pa.Table.from_pylist(records) # Write data using PyIceberg table.append(df) Dockerfile ãããŠã³ããŒãããŸããããã¯ãLambda 颿°ã®ã³ã³ããã€ã¡ãŒãžãå®çŸ©ããŸãã aws s3 cp s3://aws-blogs-artifacts-public/artifacts/BDB-5013/Dockerfile . requirements.txt ãããŠã³ããŒãããŸããããã¯ãLmbmda 颿°ã«å¿
èŠãª Python ããã±ãŒãžãå®çŸ©ããŠããŸãã aws s3 cp s3://aws-blogs-artifacts-public/artifacts/BDB-5013/requirements.txt . ãã®æç¹ã§ãäœæ¥ãã£ã¬ã¯ããªã«ã¯æ¬¡ã® 3 ã€ã®ãã¡ã€ã«ãå«ãŸããŠããã¯ãã§ãã Dockerfile lambda_function.py requirements.txt ç°å¢å€æ°ãèšå®ããŸãã ãèªåã® AWS ã¢ã«ãŠã³ã ID ã«çœ®ãæããŠãã ãã: export AWS_ACCOUNT_ID= <account_id> Docker ã€ã¡ãŒãžããã«ãããŸã: docker build --provenance=false -t localhost/pyiceberg-lambda . # Confirm built image docker images | grep pyiceberg-lambda ã€ã¡ãŒãžã«ã¿ã°ãèšå®ããŸã: docker tag localhost/pyiceberg-lambda:latest ${AWS_ACCOUNT_ID}.dkr.ecr.${AWS_REGION}.amazonaws.com/pyiceberg-lambda-repository:latest AWS CloudFormation ã«ãã£ãŠäœæããã ECR ãªããžããªã«ãã°ã€ã³ããŸã: aws ecr get-login-password --region ${AWS_REGION} | docker login --username AWS --password-stdin ${AWS_ACCOUNT_ID}.dkr.ecr.${AWS_REGION}.amazonaws.com ã€ã¡ãŒãžã ECR ãªããžããªã«ããã·ã¥ããŸã: docker push ${AWS_ACCOUNT_ID}.dkr.ecr.${AWS_REGION}.amazonaws.com/pyiceberg-lambda-repository:latest Amazon ECR ã«ããã·ã¥ããã³ã³ããã€ã¡ãŒãžã䜿çšããŠãLambda 颿°ãäœæããŸã: aws lambda create-function \ --function-name pyiceberg-lambda-function \ --package-type Image \ --code ImageUri=${AWS_ACCOUNT_ID}.dkr.ecr.${AWS_REGION}.amazonaws.com/pyiceberg-lambda-repository:latest \ --role arn:aws:iam::${AWS_ACCOUNT_ID}:role/pyiceberg-lambda-function-role-${AWS_REGION} \ --environment "Variables={ICEBERG_TABLE_NAME=nyc_yellow_table, GLUE_DATABASE_NAME=pyiceberg_lambda_blog_database}" \ --region ${AWS_REGION} \ --timeout 60 \ --memory-size 1024 次ã®ã»ã¯ã·ã§ã³ã§ç¢ºèªãããããå°ãªããšã 5 å颿°ãåŒã³åºããŠè€æ°ã®ã¹ãããã·ã§ãããäœæããŠãã ãããä»åã¯æåã§é¢æ°ãåŒã³åºããŠã€ãã³ãé§ååã®ããŒã¿åã蟌ã¿ãã·ãã¥ã¬ãŒãããŠããŸãããå®éã®ã·ããªãªã§ã¯ Lambda 颿°ãã€ãã³ãé§ååã§èªåçã«åŒã³åºãããŸãã aws lambda invoke \ --function-name arn:aws:lambda:${AWS_REGION}:${AWS_ACCOUNT_ID}:function:pyiceberg-lambda-function \ --log-type Tail \ outputfile.txt \ --query 'LogResult' | tr -d '"' | base64 -d ãããŸã§ã§ãLambda 颿°ããããã€ããŠå®è¡ããŸããããã®é¢æ°ã¯ã pyiceberg_lambda_blog_database ããŒã¿ããŒã¹å
ã« nyc_yellow_table Iceberg ããŒãã«ãäœæããŸãããŸãããã®ããŒãã«ã«ãµã³ãã«ããŒã¿ãçæããŠæ¿å
¥ããŸããåŸã®æé ã§ããã®ããŒãã«ã®ã¬ã³ãŒãã確èªããŸãã ã³ã³ããã䜿çšãã Lambda 颿°ã®æ§ç¯ã«é¢ãã詳现æ
å ±ã¯ã ã³ã³ããã€ã¡ãŒãžã䜿çšãã Lambda 颿°ã®äœæ ãã芧ãã ããã Jupyter ã䜿ã£ã PyIceberg ã«ããããŒã¿æ¢çŽ¢ ãã®ã»ã¯ã·ã§ã³ã§ã¯ãData Catalog ã«ç»é²ããã Iceberg ããŒãã«ã®ããŒã¿ã«ã¢ã¯ã»ã¹ããåæããæ¹æ³ã瀺ããŸããPyIceberg ã䜿çšãã Jupyter ããŒãããã¯ãããLambda 颿°ã«ãã£ãŠäœæãããã¿ã¯ã·ãŒéè¡ããŒã¿ã«ã¢ã¯ã»ã¹ããæ°ããã¬ã³ãŒããå°çãããã³ã«äœæãããç°ãªãã¹ãããã·ã§ãããæ€æ»ããŸãããŸããéèŠãªã¹ãããã·ã§ããã«ã¿ã°ãä»ããŠä¿æãããããªãåæã®ããã«æ°ããããŒãã«ãäœæããŸãã SageMaker AI ããŒãããã¯ã€ã³ã¹ã¿ã³ã¹ã§ Jupyter ã䜿ã£ãŠããŒãããã¯ãéãã«ã¯ãæ¬¡ã®æé ãå®è¡ããŠãã ããã SageMaker AI ã³ã³ãœãŒã«ã§ãããã²ãŒã·ã§ã³ãã€ã³ãã Notebooks ãéžæããŸãã CloudFormation ãã³ãã¬ãŒãã䜿çšããŠäœæããããŒãããã¯ã€ã³ã¹ã¿ã³ã¹ã®æšªã«ãã Open JupyterLab ãéžæããŸãã ããŒããã㯠ãããŠã³ããŒãããSageMaker AI ããŒãããã¯ã® Jupyter ç°å¢ã§éããŠãã ããã ã¢ããããŒããã pyiceberg_notebook.ipynb ãéããŸãã ã«ãŒãã«éžæã®ãã€ã¢ãã°ã§ã¯ãããã©ã«ãã®ãªãã·ã§ã³ã®ãŸãŸã«ã㊠Select ãéžæããŸãã ããããã¯ãã»ã«ãé çªã«å®è¡ããŠããŒãããã¯ãé²ããŠãããŸãã ã«ã¿ãã°ãžã®æ¥ç¶ãšããŒãã«ã®ã¹ãã£ã³ PyIceberg ã䜿çšã㊠Iceberg ããŒãã«ã«ã¢ã¯ã»ã¹ã§ããŸããæ¬¡ã®ã³ãŒãã¯ãAWS Glue Iceberg REST ãšã³ããã€ã³ãã«æ¥ç¶ãã pyiceberg_lambda_blog_database ããŒã¿ããŒã¹äžã® nyc_yellow_table ããŒãã«ãèªã¿èŸŒã¿ãŸãã import pyarrow as pa from pyiceberg.catalog import load_catalog import boto3 # Set AWS region sts = boto3.client('sts') region = sts._client_config.region_name # Configure catalog connection properties catalog_properties = { "type": "rest", "uri": f"https://glue.{region}.amazonaws.com/iceberg", "s3.region": region, "rest.sigv4-enabled": "true", "rest.signing-name": "glue", "rest.signing-region": region } # Specify database and table names database_name = "pyiceberg_lambda_blog_database" table_name = "nyc_yellow_table" # Load catalog and get table catalog = load_catalog(**catalog_properties) table = catalog.load_table(f"{database_name}.{table_name}") Iceberg ããŒãã«ãããã«ããŒã¿ã Apache Arrow ããŒãã«ãšããŠã¯ãšãªãã Pandas ã® DataFrame ã«å€æã§ããŸãã ã¹ãããã·ã§ããã®æäœ Iceberg ã®éèŠãªæ©èœã® 1 ã€ãã¹ãããã·ã§ããããŒã¹ã®ããŒãžã§ã³ç®¡çã§ããã¹ãããã·ã§ããã¯ãããŒãã«ã®ããŒã¿ã«å€æŽããããã³ã«èªåçã«äœæãããŸããæ¬¡ã®äŸã®ããã«ãç¹å®ã®ã¹ãããã·ã§ããããããŒã¿ãååŸã§ããŸãã # Get data from a specific snapshot ID snapshot_id = snapshots.to_pandas()["snapshot_id"][3] snapshot_pa_table = table.scan(snapshot_id=snapshot_id).to_arrow() snapshot_df = snapshot_pa_table.to_pandas() ã¹ãããã·ã§ããã«åºã¥ããŠãä»»æã®æç¹ã®éå»ã®ããŒã¿ãšçŸåšã®ããŒã¿ãæ¯èŒã§ããŸãããã®å Žåãææ°ã®ããŒãã«ãšã¹ãããã·ã§ããããŒãã«ã®éã®ããŒã¿ååžã®éããæ¯èŒããŠããŸãã # Compare the distribution of total_amount in the specified snapshot and the latest data. import matplotlib.pyplot as plt plt.figure(figsize=(4, 3)) df['total_amount'].hist(bins=30, density=True, label="latest", alpha=0.5) snapshot_df['total_amount'].hist(bins=30, density=True, label="snapshot", alpha=0.5) plt.title('Distribution of total_amount') plt.xlabel('total_amount') plt.ylabel('relative Frequency') plt.legend() plt.show() ã¹ãããã·ã§ããã®ã¿ã°ä»ã ç¹å®ã®ã¹ãããã·ã§ããã«ä»»æã®ååãä»ããŠã¿ã°ä»ãããåŸã§ãã®ååã§ç¹å®ã®ã¹ãããã·ã§ãããååŸã§ããŸããããã¯ãéèŠãªã€ãã³ãã®ã¹ãããã·ã§ããã管çããéã«äŸ¿å©ã§ãã ãã®äŸã§ã¯ãcheckpointTag ã¿ã°ãæå®ããŠã¹ãããã·ã§ãããžã¯ãšãªãããŠããŸããããã§ã¯ polars ã䜿çšããŠãæ¢åã®ã«ã©ã tpep_dropoff_datetime ãš tpep_pickup_datetime ã«åºã¥ã㊠trip_duration ãšããæ°ããã«ã©ã ã远å ããããšã§ãæ°ãã DataFrame ãäœæããŠããŸãã # retrive tagged snapshot table as polars data frame import polars as pl # Get snapshot id from tag name df = table.inspect.refs().to_pandas() filtered_df = df[df["name"] == tag_name] tag_snapshot_id = filtered_df["snapshot_id"].iloc[0] # Scan Table based on the snapshot id tag_pa_table = table.scan(snapshot_id=tag_snapshot_id).to_arrow() tag_df = pl.from_arrow(tag_pa_table) # Process the data adding a new column "trip_duration" from check point snapshot. def preprocess_data(df): df = df.select(["vendorid", "tpep_pickup_datetime", "tpep_dropoff_datetime", "passenger_count", "trip_distance", "fare_amount"]) df = df.with_columns( ((pl.col("tpep_dropoff_datetime") - pl.col("tpep_pickup_datetime")) .dt.total_seconds() // 60).alias("trip_duration")) return df processed_df = preprocess_data(tag_df) display(processed_df) print(processed_df["trip_duration"].describe()) åŠçæžã¿ã® DataFrame ãã trip_duration åã䜿ã£ãŠæ°ããããŒãã«ãäœæããŸãããã®ã¹ãããã¯ãå°æ¥ã®åæã®ããã«ããŒã¿ãæºåããæ¹æ³ã瀺ããŠããŸããåºã«ãªãããŒãã«ã倿ŽãããŠããŠããã¿ã°ã䜿ãããšã§ãåŠçæžã¿ããŒã¿ãåç
§ããŠããããŒã¿ã®ã¹ãããã·ã§ãããæç€ºçã«æå®ã§ããŸãã # write processed data to new iceberg table account_id = sts.get_caller_identity()["Account"] new_table_name = "processed_" + table_name location = f"s3://pyiceberg-lambda-blog-{account_id}-{region}/{database_name}/{new_table_name}" pa_new_table = processed_df.to_arrow() schema = pa_new_table.schema identifier = f"{database_name}.{new_table_name}" new_table = catalog.create_table( identifier=identifier, schema=schema, location=location ) # show new table's schema print(new_table.schema()) # insert processed data to new table new_table.append(pa_new_table) åŠçæžã¿ããŒã¿ããäœæããæ°ããããŒãã«ã Athena ã§ã¯ãšãªããIceberg ããŒãã«ã®çžäºéçšæ§ãå®èšŒããŸãããã Amazon Athena ããã®ããŒã¿ã¯ãšãª åã®ã»ã¯ã·ã§ã³ã®ããŒãããã¯ããäœæããã pyiceberg_lambda_blog_database.processed_nyc_yellow_table ããŒãã«ãã Athena ã¯ãšãªãšãã£ã¿ ã§ã¯ãšãªã§ããŸã: SELECT * FROM "pyiceberg_lambda_blog_database"."processed_nyc_yellow_table" limit 10 ; ãããã®ã¹ããããéããŠãLambda ãš AWS Glue Iceberg REST ãšã³ããã€ã³ãã䜿çšãããµãŒããŒã¬ã¹ã®ããŒã¿åŠçãœãªã¥ãŒã·ã§ã³ãæ§ç¯ãå®éã«å©çšããæµããäœéšããŸããããŸããPyIceberg ã䜿çšããŠã¹ãããã·ã§ãã管çãããŒãã«æäœãå«ã Python ã«ããããŒã¿ç®¡çãšåæãè¡ããŸãããããã«ãå¥ã®ãšã³ãžã³ã§ãã Athena ã䜿çšããŠã¯ãšãªãå®è¡ããIceberg ããŒãã«ã®äºææ§ã瀺ããŸããã ã¯ãªãŒã³ã¢ãã ãã®èšäºã§äœ¿çšãããªãœãŒã¹ãã¯ãªãŒã³ã¢ããããã«ã¯ãæ¬¡ã®æé ãå®è¡ããŠãã ããã Amazon ECR ã³ã³ãœãŒã«ã§ããªããžã㪠pyiceberg-lambda-repository ã«ç§»åãããªããžããªå
ã®ãã¹ãŠã®ã€ã¡ãŒãžãåé€ããŸãã CloudShell ã§ãäœæ¥ãã£ã¬ã¯ã㪠pyiceberg_blog ãåé€ããŸãã Amazon S3 ã³ã³ãœãŒã«ã§ãCloudFormation ãã³ãã¬ãŒãã䜿çšããŠäœæãã S3 ãã±ãã pyiceberg-lambda-blog- - ã«ç§»åãããã±ããã空ã«ããŸãã ãªããžããªãšãã±ããã空ã«ãªã£ãããšã確èªããããCloudFormation ã¹ã¿ã㯠pyiceberg-lambda-blog-stack ãåé€ããŸãã Docker ã€ã¡ãŒãžã䜿çšããŠäœæãã Lambda 颿° pyiceberg-lambda-function ãåé€ããŸãã çµè« ãã®èšäºã§ã¯ãAWS Glue Data Catalog ãš PyIceberg ã䜿çšããããšã§ãå
ç¢ãªããŒã¿ç®¡çæ©èœãç¶æããªãããå¹ççã§è»œéãªããŒã¿ã¯ãŒã¯ãããŒãå®çŸã§ããããšã瀺ããŸãããããŒã ãã€ã³ãã©ã¹ãã©ã¯ãã£ã®äŸåé¢ä¿ãæå°éã«æããªãããIceberg ã®åŒ·åãªæ©èœã掻çšã§ããããšã玹ä»ããŸããããã®ã¢ãããŒãã«ãããçµç¹ã¯åæ£ã³ã³ãã¥ãŒãã£ã³ã°ãªãœãŒã¹ã®ã»ããã¢ããã管çã®è€éããªãã«ãããã« Iceberg ããŒãã«ã®å©çšãéå§ã§ããŸãã Iceberg ã®æ©èœãäœãããŒãã«ã§å°å
¥ããããšããŠããçµç¹ã«ãšã£ãŠãããã¯ç¹ã«äŸ¡å€ããããŸããPyIceberg ã®è»œéãªæ§è³ªã«ãããããŒã ã¯ããã« Iceberg ããŒãã«ã䜿ãå§ããããšãã§ããæ
£ã芪ããã ããŒã«ã䜿çšãã远å ã®åŠç¿ãæå°éã«æããããšãã§ããŸããããŒã¿ããŒãºãé«ãŸãã°ãåã Iceberg ããŒãã«ã Athena ã AWS Glue ãªã©ã® AWS åæãµãŒãã¹ããã·ãŒã ã¬ã¹ã«ã¢ã¯ã»ã¹ã§ããå°æ¥çãªã¹ã±ãŒã©ããªãã£ãžã®æç¢ºãªãã¹ãæäŸãããŸãã PyIceberg ãš AWS ã®åæãµãŒãã¹ã®è©³çްã«ã€ããŠã¯ã PyIceberg ã®ããã¥ã¡ã³ã ãš Apache Iceberg ãšã¯? ãåç
§ããããšããå§ãããŸãã èè
ã«ã€ã㊠Sotaro Hikita ã¯ãAWS ã§ã¢ããªãã£ã¯ã¹ã«ç¹åããã¹ãã·ã£ãªã¹ããœãªã¥ãŒã·ã§ã³ã¢ãŒããã¯ãã§ãããã°ããŒã¿æè¡ãšãªãŒãã³ãœãŒã¹ãœãããŠã§ã¢ãæ±ã£ãŠããŸããä»äºä»¥å€ã§ã¯ããã€ãçŸå³ããé£ã¹ç©ãæ¢ãæ±ããæè¿ã¯ãã¶ã«å€¢äžã«ãªã£ãŠããŸãã Shuhei Fukami ã¯ãAWS ã«ãããã¢ããªãã£ã¯ã¹ã«ç¹åããã¹ãã·ã£ãªã¹ããœãªã¥ãŒã·ã§ã³ã¢ãŒããã¯ãã§ããè¶£å³ã§æçãããã®ã奜ãã§ãæè¿ã¯ãã¶äœãã«ã¯ãŸã£ãŠããŸãã