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