Delta Live Tables and Pivoting
Options
- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
07-18-2024 10:51 AM
Hello,
I'm trying to create a DLT pipeline where I read data as a streaming dataset from a Kafka source, save it in a table, and then filter, transform, and pivot the data. However, I've encountered an issue: DLT doesn't support pivoting, and using foreachBatch doesn't allow me to return a streaming query to the DLT function.
What is the best way to handle this?
Here's the code I'm trying to run:
def parse_event(df, batch_id, event_id, target):
result_df = (df.filter(col("event_id") == event_id)
.withColumn("parsed_data", from_json(col("message_content"), "map<string,string>"))
.select("time", "user", "event_id", explode(col("parsed_data")).alias("key", "value"))
.groupby("time", "user", "event_id").pivot("key").agg(first("value"))
)
(result_df.write
.format("delta")
.mode("append")
.saveAsTable(target))
@dlt.table(
name="bronze_event_xxxx",
table_properties={"quality": "bronze",
"pipelines.reset.allowed": "true"},
temporary=False)
def create_bronze_table():
df = (dlt.read_stream("bronze_event")
.writeStream
.foreachBatch(lambda df, epoch_id: parse_eventlog(df, epoch_id, "xxxx", "dev.dev.bronze_event_xxxx"))
.outputMode("append")
.start()
)
return (df)