Options
- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
06-24-2024 04:54 AM - edited 06-24-2024 04:57 AM
UPDATE: tried adding writeStream.start() like error suggested + as per other posts and ended up with following error/code:
@dlt.table(
name="ETD_Bz",
temporary=False)
def Bronze():
return (spark.readStream
.format("delta")
.option("skipChangeCommits", "true")
.table("default.tbl_raw_etd_data")
)
# function that flattens json
def process_raw_data(df, batchId) :
json_schema = spark.read.json(df.rdd.map(lambda row: row.JsonString)).schema
kafka_df = df.withColumn("JsonStruct", from_json(col("JsonString"), json_schema))
fj = FlattenJson()
kafka_flattened_json = fj.flatten_json(kafka_df)
return kafka_flattened_json
@dlt.table(
name="ETD_Flattened_Bz",
spark_conf = {"spark.databricks.delta.schema.autoMerge.enabled" : "true"},
temporary=False)
def Bronze_Flattend():
stream = (spark.readStream
.format("delta")
.option("skipChangeCommits", "true")
.table("live.ETD_Bz")
.writeStream
.format("json")
.outputMode("append")
.foreachBatch(process_raw_data)
.table("default.tbl_bz_tmp_etd_data")
.start()
.awaitTermination())
return (spark.readStream
.format("delta")
.option("skipChangeCommits", "true")
.table("default.tbl_bz_tmp_etd_data")
)
getting following error:
"py4j.protocol.Py4JJavaError: An error occurred while calling o687.toTable. : org.apache.spark.SparkClassNotFoundException: [DATA_SOURCE_NOT_FOUND] Failed to find data source: foreachBatch. Please find packages at `https://spark.apache.org/third-party-projects.html`."