You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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要求的startingOffsets JSON格式字符串
  • 将该字符串传入流式读取配置中

修改后的完整脚本

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 10:50:06