Pri*_*ava 3 apache-spark spark-structured-streaming
我有一个下沉到Kafka的结构化流查询。该查询具有复杂的聚合逻辑。
我想将此查询的输出DF下沉到多个Kafka主题,每个主题都划分在不同的“键”列上。我不想为每个不同的Kafka主题都拥有多个Kafka接收器,因为那意味着要运行多个流查询-每个Kafka主题都需要一个查询,特别是因为我的聚合逻辑很复杂。
问题:
有没有一种方法可以将结构化流查询的结果输出到多个Kafka主题,每个主题具有不同的键列,而不必执行多个流查询呢?
如果不是这样,级联多个查询以使第一个查询进行复杂的聚合并将输出写入Kafka,然后其他查询仅读取第一个查询的输出并将其主题写入Kafka会比较有效吗,从而避免了复杂的操作再次聚合?
在此先感谢您的帮助。
所以答案有点像盯着我。也有记录。下方链接。
一个查询就可以写入多个Kafka主题。如果要编写的数据框具有名为“ topic”的列(以及“ key”和“ value”列),它将把行的内容写入该行中的主题。这会自动工作。因此,您唯一需要弄清楚的是如何生成该列的值。
我也在寻找这个问题的解决方案,在我的情况下它不一定是 kafka sink。我想在 sink1 中写入数据帧的一些记录,而在 sink2 中写入一些其他记录(取决于某些条件,而不在 2 个流查询中读取相同的数据两次)。目前,根据当前的实现,它似乎不可能(DataSource.scala 中的 createSink() 方法提供对单个接收器的支持)。
然而,在 Spark 2.4.0 中有一个新的 api:foreachBatch() 它将提供数据帧微批处理的句柄,可用于缓存数据帧、写入不同的接收器或在取消缓存 aagin 之前多次处理。像这样的东西:
streamingDF.writeStream.foreachBatch { (batchDF: DataFrame, batchId: Long) =>
batchDF.cache()
batchDF.write.format(...).save(...) // location 1
batchDF.write.format(...).save(...) // location 2
batchDF.uncache()
}
Run Code Online (Sandbox Code Playgroud)
现在这个功能在数据块运行时可用:https ://docs.databricks.com/spark/latest/structured-streaming/foreach.html#reuse-existing-batch-data-sources-with-foreachbatch
编辑 15/Nov/18: 它现在在 Spark 2.4.0 中可用(https://issues.apache.org/jira/browse/SPARK-24565)
| 归档时间: |
|
| 查看次数: |
1746 次 |
| 最近记录: |