【问题标题】:How to process DynamoDB Stream in a Spark streaming application如何在 Spark 流应用程序中处理 DynamoDB Stream
【发布时间】:2017-04-16 20:18:49
【问题描述】:

我想使用 Spark Streaming 应用程序中的 DynamoDB Stream。

Spark 流式传输使用 KCL 从 Kinesis 读取数据。有一个库可以让 KCL 能够从 DynamoDB 流中读取数据:dynamodb-streams-kinesis-adapter

但是可以将这个库插入到 spark 中吗?有人做过吗?

我使用的是 Spark 2.1.0。

我的备份计划是让另一个应用程序从 DynamoDB 流中读取到 Kinesis 流中。

谢谢

【问题讨论】:

  • 到目前为止你尝试了什么?
  • 我已经能够通过调整从 DynamoDB 流中使用:KinesisUtils、KinesisInputDStream 和 KinesisReceiver。真正的变化在于我使用 com.amazonaws.services.dynamodbv2.streamsadapter.StreamsWorker 的 KinesisReceiver。

标签: apache-spark amazon-dynamodb amazon-kinesis


【解决方案1】:

实现KinesisInputDStream的方法是使用dynamodb-streams-kinesis-adapter提供的worker official guidelines 建议如下:

final Worker worker = StreamsWorkerFactory .createDynamoDbStreamsWorker( recordProcessorFactory, workerConfig, adapterClient, amazonDynamoDB, amazonCloudWatchClient);

从 Spark 的角度来看,它是在 KinesisInputDStream.scala 中的 kinesis-asl 模块下实现的

我已经为 Spark 2.4.0 尝试过这个。这是我的回购。它需要一点改进,但​​可以完成工作

https://github.com/ravi72munde/spark-dynamo-stream-asl

修改完KinesisInputDStream后,我们就可以使用了,如下图。 val stream = KinesisInputDStream.builder .streamingContext(ssc) .streamName("sample-tablename-2") .regionName("us-east-1") .initialPosition(new Latest()) .checkpointAppName("sample-app") .checkpointInterval(Milliseconds(100)) .storageLevel(StorageLevel.MEMORY_AND_DISK_2) .build()

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-08-17
    • 2017-02-24
    • 2017-01-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多