【问题标题】:Event Hub: org.apache.spark.sql.AnalysisException: Required attribute 'body' not found事件中心:org.apache.spark.sql.AnalysisException:找不到所需的属性“body”
【发布时间】:2021-09-18 02:09:53
【问题描述】:

我正在尝试将更改数据捕获写入 EventHub:

df = spark.readStream.format("delta") \
  .option("readChangeFeed", "true") \
  .option("startingVersion", 0) \
  .table("cdc_test1")

在写入 azure eventthub 时,它期望内容为 body 属性:

df.writeStream.format("eventhubs").option("checkpointLocation", checkpointLocation).outputMode("append").options(**ehConf).start()

它给出了异常

org.apache.spark.sql.AnalysisException: Required attribute 'body' not found.
    at org.apache.spark.sql.eventhubs.EventHubsWriter$.$anonfun$validateQuery$2(EventHubsWriter.scala:53)

我不确定如何将整个流包装成一个主体。我想,我需要另一个流对象,它的列体的值为“df”(原始流)作为字符串。我无法做到这一点。请帮忙!

【问题讨论】:

    标签: databricks azure-databricks azure-eventhub delta-lake


    【解决方案1】:

    您只需要使用函数struct(将所有列编码为一个对象)和类似to_json(从对象创建单个值 - 您可以使用其他函数,如@ 987654324@,或to_avro,但取决于与消费者的合同)。代码如下所示:

    df.select(F.to_json(F.struct("*")).alias("body"))\
        .writeStream.format("eventhubs")\
        .option("checkpointLocation", checkpointLocation)\
        .outputMode("append")\
        .options(**ehConf).start()
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2011-02-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-10-22
      • 1970-01-01
      相关资源
      最近更新 更多