【问题标题】:How to write multiple columns in spark dataframe to kafka queue如何将火花数据框中的多列写入kafka队列
【发布时间】:2019-05-28 09:49:30
【问题描述】:

我知道我们可以将 spark 与 kafka 集成,并将数据帧以 key 和 value 格式写入 kafka 队列,如下所示

df - 数据框

 df.withColumnRenamed("Column_1", "key")
 .withColumnRenamed("Column_2", "value")
 .write()
 .format("kafka")
 .option("kafka.bootstrap.servers", "host1:port1,host2:port2")
 .save()

但是如何将第 3、4、5 列和许多列写入 kafka 队列? 如何一次将整行写入kafka队列?

欢迎提出任何建议

【问题讨论】:

    标签: java apache-spark dataframe apache-kafka


    【解决方案1】:

    Kafka 只获取 (key, value) 形式的消息。因此,您必须将列聚合为一个值(如 JSON)。这里是例子

    这应该可行:(构造适当的value_fields

    import org.apache.spark.sql.functions._
    
    val value_fields = df.columns.filter(_ != "Column_1") 
    
    df
    .withColumnRenamed("Column_1", "key")
    .withColumn("value", to_json(struct(value_fields.map(col(_)):_*)))
    .select("key", "value")
    .write()
    .format("kafka")
    .option("kafka.bootstrap.servers", "host1:port1,host2:port2")
    .save()
    

    【讨论】:

    • 你能说得更具体些吗
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-10-13
    • 1970-01-01
    • 1970-01-01
    • 2019-08-30
    • 1970-01-01
    • 2018-11-30
    • 1970-01-01
    相关资源
    最近更新 更多