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

Kinesis Data Viewer无法查看数据及Glue DataFrame读取问题排查

Kinesis数据流读取问题排查

我已经向Kinesis数据流中插入了数据,通过序列号能查到记录,但选择Latest选项却看不到数据。AWS文档里对Latest的描述让我困惑:

显示分片最新记录之后的记录,以便始终读取分片的最新数据

“最新记录之后怎么会有数据?”而且我插入的最新数据也看不到。另外,按照AWS官方博客教程操作时,Trim Horizon选项也无法显示数据。

问题1:Latest选项无法加载数据?是否需要修改putRecord API调用?

当前调用代码:

response = kinesis_client.put_record(StreamARN=SECURITY_LAKE_AZURE_STREAM_ARN,
            Data=json.dumps(record),
            PartitionKey="time"
            )

问题2:如何配置Glue DataFrame的连接选项以读取这些数据?

设置"startingPosition": "earliest"无法获取任何数据,当前DataFrame代码:

dataframe_KinesisStream_node1 = glueContext.create_data_frame.from_options(
    connection_type="kinesis",
    connection_options={
        "typeOfData": "kinesis",
        "streamARN": SECURITY_LAKE_AZURE_STREAM_ARN,
        "classification": "json",
        "startingPosition": "earliest",
        "inferSchema": "true",
    },
    transformation_ctx="dataframe_KinesisStream_node1",
)

问题1解答

关于Latest选项的理解纠正

AWS文档的描述容易产生歧义,正确逻辑是:选择Latest作为读取位置时,Kinesis会从你发起读取请求的当前时刻之后产生的新记录开始返回数据。也就是说,你之前插入的历史记录不会被包含在内——它只监听“未来”的新数据,而非返回已有的最新记录,这就是你看不到已插入数据的核心原因。

putRecord调用的检查点

你的代码本身无语法问题,但需要确认以下几点:

  • 验证PartitionKey路由:如果"time"是固定字符串,所有记录会被路由到同一个分片。需通过putRecord返回的ShardId确认你在控制台查看的分片是实际接收数据的分片。
  • 确认写入成功:用putRecord返回的SequenceNumber,在控制台指定对应分片和序列号,确认记录确实存在。
  • 检查数据流状态:确保数据流处于ACTIVE状态,未被暂停或删除。

问题2解答

Glue DataFrame读取失败的常见原因及修复

设置"startingPosition": "earliest"仍读不到数据,可从以下方向排查:

  1. 权限缺失:Glue作业执行角色必须具备kinesis:DescribeStream、kinesis:GetShardIterator、kinesis:GetRecords权限,同时信任策略需允许Glue服务调用这些权限。
  2. 数据已过期:Kinesis数据流默认数据保留期为24小时,最长365天。如果数据写入时间超过保留期,earliest也读取不到。
  3. 连接参数遗漏:部分Glue版本需同时指定streamName(可从ARN末尾提取,格式为arn:aws:kinesis:region:account-id:stream/stream-name),仅用ARN可能无法正常识别。
  4. 数据格式问题:若JSON数据存在语法错误,classification: "json"和inferSchema: "true"会导致数据被过滤。建议先关闭inferSchema,以原始字符串格式读取验证数据是否存在。

修改后的Glue DataFrame配置示例

dataframe_KinesisStream_node1 = glueContext.create_data_frame.from_options(
    connection_type="kinesis",
    connection_options={
        "typeOfData": "kinesis",
        "streamARN": SECURITY_LAKE_AZURE_STREAM_ARN,
        "streamName": "your-stream-name",  # 从ARN中提取的数据流名称
        "classification": "json",
        "startingPosition": "earliest",
        "inferSchema": "false",  # 先关闭自动推断,验证数据可读性
        "startingPositionTimestamp": "2024-01-01T00:00:00Z"  # 可选,指定起始时间覆盖数据写入时段
    },
    transformation_ctx="dataframe_KinesisStream_node1",
)

额外排查步骤

  • 查看Glue作业日志:在CloudWatch中检查是否有Kinesis相关错误(如权限不足、分片不存在)。
  • 用CLI验证:在本地或EC2执行aws kinesis get-records命令,指定earliest类型的shard iterator,确认数据是否能被读取,排除Glue配置问题。

内容的提问来源于stack exchange,提问作者manu muraleedharan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 06:23:27