在Spark实时流中刷新Dataframe而不停止进程

shi*_*455 3 amazon-s3 apache-spark spark-streaming spark-dataframe snappydata

在我的应用程序中,我从Kafka队列获得了一个帐户流(使用带有kafka的Spark流)

我需要从S3获取与这些帐户相关的属性,因此我计划缓存S3结果数据帧,因为S3数据暂时不会更新至今一天,未来可能会很快变为1小时或10分钟.所以问题是如何定期刷新缓存的数据框而不停止进程.

**更新:我计划在S3中有更新时使用SNS和AWS lambda将事件发布到kafka,我的流应用程序将订阅事件并根据此事件刷新缓存的数据帧(基本上是unpersist()缓存和从S3重新加载)这是一个好方法吗?

pla*_*bre 5

最近在Spark邮件列表中询问了这个问题

据我所知,你要问的唯一方法是在新数据到达时从S3重新加载DataFrame,这意味着你必须重新创建流式DF并重新启动查询.这是因为DataFrames基本上是不可变的.

如果要在不重新加载DataFrame的情况下更新(mutate)DataFrame中的数据,则需要尝试与Spark集成或连接到Spark并允许突变的数据存储之一.我所知道的是SnappyData.