You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何读取Struct类型DataFrame列并将键值对添加至Kafka Headers?

解决方案

你只需要先从Row中取出Headers这个Struct对应的子Row,再从这个子Row里读取具体的字段值即可,修改后的代码如下:

df.foreachPartition {
  partitionData: Iterator[Row] =>
    partitionData.foreach { row =>
      val kafkaProps = getKafkaConfig()
      val producer = JsonObjectProducer.getOrCreateProducer(kafkaProps)
      val data = new ProducerRecord[String, String](kafkaTopic, row.getAs[String]("Payload"))
      if(headerFlag == "True"){
        // 先获取Struct类型的Headers对应的Row
        val headersRow = row.getAs[Row]("Headers")
        // 从子Row中读取各个字段并添加到Kafka Headers
        data.headers().add("Id", headersRow.getAs[String]("Id").getBytes(StandardCharsets.UTF_8))
        data.headers().add("key", headersRow.getAs[String]("key").getBytes(StandardCharsets.UTF_8))
        data.headers().add("ver", headersRow.getAs[String]("ver").getBytes(StandardCharsets.UTF_8))
      }
      producer.send(data, new Callback() {
        def onCompletion(rm: RecordMetadata, e: Exception) {
          if (e != null) {
            e.printStackTrace()
          }
        }
      })
    }
}

关键说明

  • 在Spark中,Struct类型的列会被封装成Row对象,所以需要先用row.getAs[Row]("Headers")取出整个Struct结构
  • 之后就可以像操作普通Row一样,用getAs[String]从这个子Row里获取对应的字段值,再转成字节数组添加到Kafka Headers中

内容的提问来源于stack exchange,提问作者abhimanyu

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.01 07:10:14