AWS Glue Streaming ETL用AvailableNow触发Kinesis时无记录处理且Checkpoint异常
问题:AWS Glue Streaming ETL用AvailableNow触发器读取Kinesis无记录处理,Checkpoint参数不一致
现象
- 用
trigger(availableNow=True)触发Glue流任务读取Kinesis Data Streams,全程无任何记录被处理 - 查看S3上的checkpoint文件,发现其中
startingPosition和iteratorType相关数据存在不一致情况
涉事代码
streaming_df: DataFrame = spark \ .readStream \ .format("kinesis") \ .option("streamName", "oxg-cdp-cdc-stream-sbx") \ .option("endpointUrl", "https://kinesis.eu-central-1.amazonaws.com") \ .option("region", "eu-central-1")\ .option("startingPosition", "earliest")\ .load() streaming_df \ .writeStream \ .option("checkpointLocation", "<s3_location>") \ .trigger(availableNow=True) \ .foreachBatch(for_each_batch_funtion) \ .start() \ .awaitTermination()
分析与解决办法
兼容性问题
目前AWS官方文档未明确标注,但实际使用中AvailableNow模式(一次性处理完所有可用数据即停止)与Kinesis连接器的checkpoint位置追踪逻辑存在适配冲突,导致checkpoint中存储的位置参数混乱,无法正确读取流数据。替代方案
- 改用
trigger(processingTime='X seconds')的连续触发模式,这是Glue与Kinesis流兼容的标准触发方式,能保证checkpoint正确记录读取位置,正常处理数据。 - 如果需要一次性处理历史数据的场景:
- 先编写批处理任务,通过
startingPosition='earliest'读取Kinesis所有历史数据 - 再切换为流模式处理新产生的数据,利用checkpoint衔接读取位置
- 先编写批处理任务,通过
- 改用
清理异常Checkpoint
若之前的checkpoint已存储异常数据,先删除S3上对应的checkpoint目录,再重启任务,避免旧的错误参数干扰新任务执行。
内容的提问来源于stack exchange,提问作者Alexandre Silva
相关产品推荐
相关产品推荐

