避免多个流查询

Pri*_*ava 3 apache-spark spark-structured-streaming

我有一个下沉到Kafka的结构化流查询。该查询具有复杂的聚合逻辑。

我想将此查询的输出DF下沉到多个Kafka主题,每个主题都划分在不同的“键”列上。我不想为每个不同的Kafka主题都拥有多个Kafka接收器,因为那意味着要运行多个流查询-每个Kafka主题都需要一个查询,特别是因为我的聚合逻辑很复杂。

问题:

  1. 有没有一种方法可以将结构化流查询的结果输出到多个Kafka主题,每个主题具有不同的键列,而不必执行多个流查询呢?

  2. 如果不是这样,级联多个查询以使第一个查询进行复杂的聚合并将输出写入Kafka,然后其他查询仅读取第一个查询的输出并将其主题写入Kafka会比较有效吗,从而避免了复杂的操作再次聚合?

在此先感谢您的帮助。

Pri*_*ava 6

所以答案有点像盯着我。也有记录。下方链接。

一个查询就可以写入多个Kafka主题。如果要编写的数据框具有名为“ topic”的列(以及“ key”和“ value”列),它将把行的内容写入该行中的主题。这会自动工作。因此,您唯一需要弄清楚的是如何生成该列的值。

这是记录在案的-https://spark.apache.org/docs/latest/structured-streaming-kafka-integration.html#writing-data-to-kafka


cha*_*ash 5

我也在寻找这个问题的解决方案,在我的情况下它不一定是 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