【问题标题】:How to get Kafka header's value to Spark Dataset as a single column?如何将 Kafka 标头的值作为单列获取到 Spark 数据集?
【发布时间】:2022-02-14 02:09:13
【问题描述】:

我们有一个启用了标头的 Kafka 流

  .option("includeHeaders", true)

从而使它们存储为高级数据集的列,承载带有键和值的内部结构数组:

root
 |-- topic: string (nullable = true)
 |-- key: string (nullable = true)
 |-- value: string (nullable = true)
 |-- timestamp: string (nullable = true)
 |-- headers: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- key: string (nullable = true)
 |    |    |-- value: binary (nullable = true)

我可以使用数组中的顺序访问所需的标题:

val controlDataFrame = spark
      .readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", kafkaLocation)
      .option("includeHeaders", true)
      .option("failOnDataLoss", value = false)
      .option("subscribe", "mytopic")
      .load()
      .withColumn("acceptTimestamp", element_at(col("headers"),1))
      .withColumn("acceptTimestamp2", col("acceptTimestamp.value").cast("STRING"))
 

但是这个解决方案看起来很脆弱,因为在另一端产生的标题的顺序总是可以随着更新而改变,而只有键名在那里看起来很稳定。如何通过结构键查找并提取所需的结构而不是指向数组索引?

更新。

感谢 Alex Ott 的 davice,我找到了将我想要的内容放入以下列的解决方案:

.withColumn("headers1", map_from_entries(col("headers")))
.withColumn("acceptTimestamp2", col("headers1.acceptTimestamp").cast("STRING"))

【问题讨论】:

    标签: scala apache-spark apache-kafka spark-structured-streaming


    【解决方案1】:

    您可以使用map_from_entries 函数将结构数组转换为可以按名称访问条目的映射。

    import org.apache.spark.sql.functions.map_from_entries
    
    ....
    select(map_from_entries("headers").alias("headers"), ...)
    

    但我记得,标头名称可能不是唯一的,这是将它们作为键/值对数组发送的主要原因。

    另一种方法是使用 filter 函数按名称查找标头 - 这将允许处理非唯一标头。

    附:我使用 Python 文档是因为我可以链接各个函数 - 在 Scala 文档中这并不容易。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-07-30
      • 2020-03-11
      • 2017-12-06
      • 2015-11-05
      • 1970-01-01
      • 2023-01-25
      • 2018-10-30
      相关资源
      最近更新 更多