Spark Streaming无法读取Kinesis流,作业未接收到记录求助
排查Spark作业无法接收Kinesis流数据的步骤
遇到这种情况确实挺头疼的,既然已经确认凭据没问题,那我们从几个常见的核心方向一步步排查:
1. 核对流名称与分片关联
- 再仔细确认Spark作业中配置的Kinesis流名称完全一致,包括大小写、特殊字符,有时候哪怕一个字母写错都可能导致订阅错流。
- 登录Kinesis控制台查看目标流的分片状态,同时检查Spark作业的driver日志,搜索是否有类似
Received X shards from stream [流名称]的日志,确认作业是否成功关联到了流的所有分片。如果日志里显示分片数为0,那说明作业根本没连上目标流。
2. 检查消费位置配置
这是最容易踩坑的点之一:
- 如果是首次启动的作业,确认
InitialPositionInStream参数设置:- 如果设为
LATEST,作业只会消费启动之后流入的新数据,之前已存在的旧数据不会被拉取;如果你的测试数据是启动前流入的,改成TRIM_HORIZON就能从最早的记录开始消费。
- 如果设为
- 如果是重启的作业,检查是否启用了checkpoint:
- 如果checkpoint记录的消费位置已经超过了Kinesis流的保留期(默认24小时,最长可设365天),那该位置之后的记录已经被删除,作业会一直等待新数据流入。可以尝试临时禁用checkpoint,用
TRIM_HORIZON重新启动测试。
- 如果checkpoint记录的消费位置已经超过了Kinesis流的保留期(默认24小时,最长可设365天),那该位置之后的记录已经被删除,作业会一直等待新数据流入。可以尝试临时禁用checkpoint,用
3. 验证数据格式与解析逻辑
有时候数据确实到了,但解析失败被丢弃:
- 在Spark作业中临时添加一段调试代码,直接打印Kinesis记录的原始字节内容,不要做任何解析:
运行后看控制台是否有输出,如果有,说明数据已经被接收,问题出在后续的解析逻辑上;如果还是没有输出,那问题在数据拉取环节。kinesisStream.map(record => new String(record.getData())) .foreachRDD(rdd => rdd.foreach(println)) - 确认Kinesis流入的数据编码格式(比如UTF-8)和Spark作业中使用的解码格式一致,避免因编码错误导致解析失败。
4. 排查作业日志与运行状态
- 查看Spark driver和executor的完整日志,搜索关键词
Kinesis、fetch、record,看看有没有报错信息,比如:- 网络超时错误(说明集群和Kinesis服务的网络连通性有问题)
- API调用超限(Kinesis有请求频次限制,可查看控制台的监控指标)
- 分片关闭或过期的提示
- 确认Spark Streaming上下文是否正确启动:作业中必须调用
streamingContext.start()和streamingContext.awaitTermination(),缺少任何一个都会导致作业无法开始消费。
5. 其他细节检查
- 确认Kinesis流的状态是ACTIVE,没有被暂停或删除。
- 检查Spark版本与Kinesis连接器版本是否兼容:比如Spark 3.x需要使用适配3.x版本的
spark-streaming-kinesis-asl依赖,版本不匹配可能会出现隐藏的兼容性问题。 - 如果集群在VPC内,确认是否配置了Kinesis的VPC端点,或者集群有访问公网的权限,保证能正常调用Kinesis的API。
你可以先从这些方向逐一排查,把对应的结果反馈出来,能更精准地定位问题。
内容的提问来源于stack exchange,提问作者grepIt
相关产品推荐
相关产品推荐

