【发布时间】: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