Glue脚本基于Timestamp读取Kafka事件的问题及解决方案咨询
基于时间戳从Kafka读取Glue流式数据的解决方案
问题说明
使用Glue的Spark Structured Streaming读取Kafka事件时,默认从最早offset开始消费,需要改为从指定时间戳开始读取。直接设置startingOffsets为timestamp并指定startingTimestamp的方式无效,任务无数据输出就停止。
可行方案
Spark Structured Streaming本身不支持直接通过时间戳设置起始offset,但可以通过先将时间戳转换为对应分区的offset,再指定起始offset的方式实现。具体逻辑:
- 将目标时间戳转为Kafka使用的毫秒级时间戳
- 通过Kafka AdminClient查询目标Topic每个分区中,大于等于目标时间戳的最早offset
- 构造符合Spark要求的
startingOffsetsJSON格式字符串 - 将该字符串传入流式读取配置中
修改后的完整脚本
import sys import json import traceback from datetime import datetime from kafka import KafkaAdminClient from kafka.errors import KafkaError import pyspark from pyspark import SparkContext from pyspark.sql import SparkSession from pyspark.sql.functions import * # 初始化Spark上下文 sc = SparkContext() sc.setSystemProperty("com.amazonaws.services.s3.enableV4", "true") hadoopConf = sc._jsc.hadoopConfiguration() hadoopConf.set("fs.s3a.aws.credentials.provider", "com.amazonaws.auth.profile.ProfileCredentialsProvider") hadoopConf.set("com.amazonaws.services.s3a.enableV4", "true") hadoopConf.set("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") spark = SparkSession(sc).builder.getOrCreate() def get_starting_offsets(kafka_bootstrap_servers, topic_name, target_timestamp_str): """根据目标时间戳获取Kafka Topic各分区的起始offset""" # 将ISO格式时间转为毫秒级时间戳 target_dt = datetime.fromisoformat(target_timestamp_str.replace("Z", "+00:00")) target_timestamp_ms = int(target_dt.timestamp() * 1000) # 初始化Kafka AdminClient,认证信息与流式读取保持一致 admin_client = KafkaAdminClient( bootstrap_servers=kafka_bootstrap_servers, security_protocol="SASL_SSL", sasl_mechanism="PLAIN", sasl_plain_username="USERNAME", sasl_plain_password="PASSWORD" ) try: # 获取Topic所有分区ID topic_info = admin_client.describe_topics([topic_name])[0] partition_ids = [p["partition"] for p in topic_info["partitions"]] # 逐个分区查询对应时间戳的offset offsets = {} for partition in partition_ids: offset_response = admin_client.offsets_for_times({(topic_name, partition): target_timestamp_ms}) # 若时间戳后无数据,使用分区最新offset避免任务空跑停止 offset = offset_response[(topic_name, partition)].offset if offset_response[(topic_name, partition)] else admin_client.get_watermark_offsets((topic_name, partition))["high"] offsets[str(partition)] = offset # 构造Spark要求的startingOffsets JSON格式 return json.dumps({topic_name: offsets}) except KafkaError as e: print(f"Kafka查询错误: {str(e)}") raise finally: admin_client.close() try: # 核心配置参数 KAFKA_BOOTSTRAP_SERVERS = "kafka_server_url:9092" TOPIC_NAME = "topic_name" TARGET_TIMESTAMP = "2023-06-20T00:00:00Z" # 获取基于时间戳的起始offset starting_offsets = get_starting_offsets(KAFKA_BOOTSTRAP_SERVERS, TOPIC_NAME, TARGET_TIMESTAMP) # 构建Kafka流式读取配置 options = { "kafka.sasl.jaas.config": 'org.apache.kafka.common.security.plain.PlainLoginModule required username="USERNAME" password="PASSWORD";', "kafka.sasl.mechanism": "PLAIN", "kafka.security.protocol": "SASL_SSL", "kafka.bootstrap.servers": KAFKA_BOOTSTRAP_SERVERS, "subscribe": TOPIC_NAME, "startingOffsets": starting_offsets } # 读取Kafka流式数据 df = spark.readStream.format("kafka").options(**options).load() # 转换消息格式 df = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") # 输出到S3(注意checkpoint与数据路径分离) df.writeStream.format("json") \ .option("checkpointLocation", "s3://mybucket/test/checkpoint/") \ .outputMode("append") \ .option("path", "s3://mybucket/test/data/") \ .start() \ .awaitTermination() except Exception as e: print(f"执行错误: {str(e)}") traceback.print_exc()
注意事项
- 依赖配置:Glue作业需添加
kafka-python库作为额外依赖,可在作业配置的「Python库路径」中上传对应whl包。 - 路径规范:checkpoint路径需与数据路径分开,避免数据覆盖或读取异常。
- 异常兼容:若目标时间戳后无数据,脚本自动切换为读取最新offset,保证任务持续运行等待新数据。
内容的提问来源于stack exchange,提问作者Smaillns
相关产品推荐
相关产品推荐

