You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Databricks Delta Live Tables流水线流处理异常求助

Databricks DLT流水线流处理模式异常问题

问题背景

在Databricks Workflows的Delta Live Tables中开发流水线,流处理环节不符合预期:

预期效果

  • Bronze表通过AutoLoader(cloudFiles)以spark.readStream流模式读取JSON文件
  • Silver表以dlt.read_stream流模式读取并扁平化Bronze表数据

实际表现

  • 读取数百个文件时:流水线启动后,Bronze阶段的行数/文件数需等到阶段完成才更新;Silver阶段启动后数据指标始终不更新,最终因内存错误终止
  • 读取少量文件时:流水线可正常运行结束,但各阶段的行数/文件数同样需等到阶段完成才更新

判断流水线未真正以流模式运行,寻求配置或运行逻辑上的疏漏排查。

流水线配置

{
"id": "xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx",
"clusters": [
    {
        "label": "default",
        "aws_attributes": {
            "instance_profile_arn": "arn:aws:iam::xxxxxxxxxxxx:instance-profile/iam_role_example"
        },
        "autoscale": {
            "min_workers": 1,
            "max_workers": 10,
            "mode": "LEGACY"
        }
    }
],
"development": true,
"continuous": false,
"channel": "CURRENT",
"edition": "PRO",
"photon": false,
"libraries": [
    {
        "notebook": {
            "path": "/Repos/user_example@xxxxxx.xx/dms/bronze_job"
        }
    }
],
"name": "01-landing-task-1",
"storage": "dbfs:/pipelines/xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx",
"configuration": {
    "SCHEMA": "example_schema",
    "RAW_MOUNT_NAME": "xxxx",
    "DELTA_MOUNT_NAME": "xxxx",
    "spark.sql.parquet.enableVectorizedReader": "false"
},
"target": "landing"}

流水线代码

Silver表实际包含约30个get_json_object列:

import dlt
import pyspark.sql.functions as F
import pyspark.sql.types as T
from pyspark.sql.window import Window

RAW_MOUNT_NAME = spark.conf.get("RAW_MOUNT_NAME")
SCHEMA = spark.conf.get("SCHEMA")
SOURCE = spark.conf.get("SOURCE")
TABLE_NAME = spark.conf.get("TABLE_NAME")
PRIMARY_KEY_PATH = spark.conf.get("PRIMARY_KEY_PATH")

@dlt.table(
    name=f"{SCHEMA}_{TABLE_NAME}_bronze",
    table_properties={
        "quality": "bronze"
    }
)
def bronze_job():
    load_path = f"/mnt/{RAW_MOUNT_NAME}/{SOURCE}/5e*"

    return spark \
            .readStream \
            .format("text") \
            .option("encoding", "UTF-8") \
            .load(load_path) \
            .select("value", "_metadata") \
            .withColumnRenamed("value", "json") \
            .withColumn("id", F.expr(f"get_json_object(json, '$.{PRIMARY_KEY_PATH}')")) \
            .withColumn("_etl_timestamp", F.col("_metadata.file_modification_time")) \
            .withColumn("_metadata", F.col("_metadata").cast(T.StringType())) \
            .withColumn("_etl_operation", F.lit("U")) \
            .withColumn("_etl_to_delete", F.lit(False)) \
            .withColumn("_etl_file_name", F.input_file_name()) \
            .withColumn("_etl_job_processing_timestamp", F.current_timestamp()) \
            .withColumn("_etl_table", F.lit(f"{TABLE_NAME}")) \
            .withColumn("_etl_partition_date", F.to_date(F.col("_etl_timestamp"), "yyyy-MM-dd")) \
            .select("_etl_operation", "_etl_timestamp", "id", "json", "_etl_file_name", "_etl_job_processing_timestamp", "_etl_table", "_etl_partition_date", "_etl_to_delete", "_metadata")

@dlt.table(
name=f"{SCHEMA}_{TABLE_NAME}_silver",
table_properties = {
        "quality": "silver",
        "delta.autoOptimize.optimizeWrite": "true", 
        "delta.autoOptimize.autoCompact": "true"
}
)
def silver_job():
    df = dlt.read_stream(f"{SCHEMA}_{TABLE_NAME}_bronze").where("_etl_table == 'extraction'")
    return df.select(
                    df.id.alias('medium_id'),
                    F.get_json_object(df.json, '$.request').alias('request_id'))

问题排查与解决方案

1. 运行模式配置问题

流水线配置中"continuous": false,当前为触发式运行(Triggered),这种模式会一次性处理所有可用数据,处理完成后就终止,表现类似批处理。如果要实现持续流处理,需要将其改为"continuous": true,启用持续运行模式。

2. Bronze表未使用AutoLoader的cloudFiles格式

代码中Bronze表用的是format("text"),而非AutoLoader的cloudFiles格式。AutoLoader需要指定format("cloudFiles")并配置相关参数,示例修改:

return spark \
        .readStream \
        .format("cloudFiles") \
        .option("cloudFiles.format", "text") \
        .option("encoding", "UTF-8") \
        .option("cloudFiles.schemaLocation", f"/mnt/{DELTA_MOUNT_NAME}/schema/{SCHEMA}_{TABLE_NAME}")  # 用于自动推断/保存 schema
        .load(load_path)

使用cloudFiles格式才会启用AutoLoader的增量监听、自动Schema进化等流处理特性,否则普通的readStream + text格式会一次性读取匹配路径下的所有文件,表现为批处理行为。

3. 内存问题优化

  • 当处理大量文件时,可添加流处理触发间隔,控制每次处理的数据量:
    在Bronze表的readStream后添加.trigger(processingTime='5 minutes')(根据实际情况调整),避免一次性加载过多数据导致内存溢出。
  • Silver表中大量使用get_json_object会增加计算开销,建议先将JSON字符串解析为结构化数据,再提取字段,减少重复解析:
    df = dlt.read_stream(f"{SCHEMA}_{TABLE_NAME}_bronze") \
            .where("_etl_table == 'extraction'") \
            .withColumn("parsed_json", F.from_json(F.col("json"), T.StructType([
                T.StructField("request", T.StringType()),
                # 其他字段根据实际JSON结构定义
            ])))
    return df.select(
            df.id.alias('medium_id'),
            F.col("parsed_json.request").alias('request_id')
            # 其他字段从parsed_json中提取
    )
    
    提前解析JSON可减少重复调用get_json_object带来的性能消耗,降低内存压力。

4. DLT流处理状态验证

启用持续运行模式后,可通过DLT流水线的“流进度”面板查看每个阶段的微批处理情况,确认是否有持续的微批生成,验证流模式是否正常运行。


内容的提问来源于stack exchange,提问作者whitehummingbird

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 06:57:11