【问题标题】:What is the best way to perform multiple filter operations on spark streaming dataframe read from Kafka?对从 Kafka 读取的火花流数据帧执行多个过滤器操作的最佳方法是什么?
【发布时间】:2021-07-22 04:19:18
【问题描述】:

我需要在从 Kafka 主题读取的 DataFrame 上应用多个过滤器,并将每个过滤器的输出发布到外部系统(如另一个 Kafka 主题)。

我读过这样的kafkaDF

val kafkaDF: DataFrame = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "localhost:9092")
  .option("subscribe", "try.kafka.stream")
  .load()
  .select(col("topic"), expr("cast(value as string) as message"))
  .filter(col("message").isNotNull &&  col("message") =!= "")
  .select(from_json(col("message"), eventsSchema).as("eventData"))
  .select("eventData.*")

我可以在此 Dataframe 上运行 foreachBatch,然后遍历过滤器列表以获取过滤后的数据,然后可以将其发布到 kafka 主题,如下所示

kafkaDF.writeStream
  .foreachBatch { (batch: DataFrame, _: Long) =>
    // List of filters that needs to be applied
    filterList.par.foreach(filterString => {
      val filteredDF = batch.filter(filterString)
      // Add some columns. 
      // Do some operations based on different filter
      filteredDF.toJSON.foreach(value => {
        // Publish a message to Kafka 
      })
    })
  }
  .trigger(Trigger.ProcessingTime("60 seconds"))
  .start()
  .awaitTermination()

但是,考虑到这么多的迭代,我不确定这是否是最好的方法。有没有比这样更好的方法?

【问题讨论】:

    标签: apache-spark apache-spark-sql spark-structured-streaming spark-kafka-integration


    【解决方案1】:

    如果您计划将来自一个 Kafka 主题的数据写入多个 Kafka 主题,您可以在写入 Kafka 时在一个单个 Dataframe 中创建一个名为“topic”的列。然后,此列中的值定义将在其中生成记录的主题。这允许您根据需要写入尽可能多的不同 Kafka 主题。

    因此,我只会将您的过滤器逻辑用作何时/否则条件,或者如果更复杂,则用作 UDF。

    下面是一个示例代码,可以帮助您入门。根据消费的 Kafka 消息的value,在filteredDf 中创建一个名为“topic”的列。如果 value = 1 则 Dataframe 记录生成到名为“out1”的主题中,否则记录生成到名为“out2”的主题中。

    val inputDf = spark.readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", "localhost:9092")
      .option("subscribe", "try.kafka.stream")
      .option("failOnDataLoss", "false")
      .load()
      .selectExpr("CAST(key AS STRING) as key", "CAST(value AS STRING) as value", "partition", "offset", "timestamp")
    
    val filteredDf = inputDf.withColumn("topic", when(filter, lit("out1")).otherwise(lit("out2")))
    
    
    val query = filteredDf
      .select(
        col("key"),
        to_json(struct(col("*"))).alias("value"),
        col("topic"))
      .writeStream
      .format("kafka")
      .option("kafka.bootstrap.servers", "localhost:9092")
      .option("checkpointLocation", "/home/michael/sparkCheckpoint/1/")
      .start()
    
    query.awaitTermination()
    

    编辑:(我一开始可能误解了你的问题)

    如果您只是想找到一种从filterList 中应用多个过滤器的好方法,您可以使用foldLeft 组合它们:

    val filter1 = col("value") === 1
    val filter2 = col("key") === 1
    val filterList = List(filter1, filter2)
    val filterAll = filterList.tail.foldLeft(filterList.head)((f1, f2) => f1.and(f2))
    
    println(filterAll)
    ((value = 1) AND (key = 1))
    

    然后将.filter(filterAll) 应用于您的数据框。

    【讨论】:

    • 谢谢!如果我可以加入过滤器,这是有道理的。但是,在这种情况下,我需要应用过滤器并基于此执行不同的操作。例如考虑我正在阅读来自 Kafka 的日志消息,我的过滤器是查找错误、警告、信息、调试等消息,并为每种类型的消息创建一个重要性得分列,然后将它们推送到 Kafka 中的单独主题
    猜你喜欢
    • 2016-11-12
    • 1970-01-01
    • 2016-06-23
    • 1970-01-01
    • 2017-04-05
    • 2020-04-16
    • 1970-01-01
    • 2023-03-18
    • 2012-10-02
    相关资源
    最近更新 更多