如何使用Spark读取Amazon SQS队列数据
Spark读取Amazon SQS的方案说明
现有连接器情况
Apache Spark官方没有内置面向Amazon SQS的原生数据源,无法像对接Kafka、Kinesis一样通过官方提供的API直接读取SQS数据生成DataFrame或DStream。
市面上存在少量第三方开源的Spark SQS连接器,基本都是个人或小团队维护,没有经过大规模生产环境验证,生产使用需要做充足的兼容性、稳定性以及性能测试。
主流实现方案
如果没有找到适配你当前Spark版本的可靠第三方连接器,可以按照你的使用场景选择以下两种自定义实现方案:
方案1:自定义Structured Streaming数据源(Spark 2.x及以上版本推荐)
如果使用Structured Streaming API,通过实现自定义Source接口完成对接:
- 实现
org.apache.spark.sql.sources.v2.reader.streaming.Source接口,核心重写getOffset、getBatch、stop三个方法 - 在
getBatch方法中调用Amazon SQS Java SDK拉取队列消息,解析后转换为InternalRow生成对应的DataFrame - 自行实现offset管理逻辑,根据业务需要配置至少一次/恰好一次消费语义,消费成功后调用SQS删除消息接口清理已完成处理的消息
你现有通过Java SDK拉取消息生成DataFrame的逻辑可以直接复用,只需要封装到自定义Source的拉取逻辑中即可,不需要完全重构现有代码。
方案2:自定义Spark Streaming Receiver(适配旧版DStream API)
如果使用的是旧版Spark Streaming的DStream API,通过实现自定义Receiver完成对接:
- 继承
org.apache.spark.streaming.receiver.Receiver类,重写onStart、onStop方法 - 在
onStart方法中启动独立的后台线程循环拉取SQS消息,调用store方法将消息存入Spark执行内存 - 配合Receiver的可靠性配置,实现对应级别的消息投递语义
内容的提问来源于stack exchange,提问作者Rajesh Kumar Dash
相关产品推荐
相关产品推荐

