Options
- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
05-15-2024 07:01 AM
That's good news. So would this be the correct sort of set up, or should I be creating all the streams first before writing the streams to a table?
customer_list = ['customer1', 'customer2', 'customer3', ...]
table_name = "bronze.customer_data_table"
for customer in customer_list:
file_path = f"wasbs://{customer}@conatiner.blob.core.windows.net/*/*.csv"
checkpoint_ path = f"/tmp/checkpoints/{customer}/_checkpoints"
cloudFile = {
"cloudFiles.format": "csv",
"cloudFiles.backfillInterval": "1 day",
"cloudFiles.schemaLocation": checkpoint_path,
"cloudFiles.schemaEvolutionMode": "rescue",
}
df = (
spark.readStream
.format("cloudFiles")
.options(**cloudFile)
.load(file_path)
)
streamQuery = (
df.writeStream.format("delta")
.option("outputMode", "append")
.option("checkpointLocation", checkpoint_path)
.trigger(once=True)
.toTable(table_name)
)