TimB
New Contributor III

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)
)