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"仍读不到数据,可从以下方向排查:
- 权限缺失:Glue作业执行角色必须具备
kinesis:DescribeStream、kinesis:GetShardIterator、kinesis:GetRecords权限,同时信任策略需允许Glue服务调用这些权限。 - 数据已过期:Kinesis数据流默认数据保留期为24小时,最长365天。如果数据写入时间超过保留期,
earliest也读取不到。 - 连接参数遗漏:部分Glue版本需同时指定
streamName(可从ARN末尾提取,格式为arn:aws:kinesis:region:account-id:stream/stream-name),仅用ARN可能无法正常识别。 - 数据格式问题:若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
相关产品推荐
相关产品推荐

