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

Databricks Spark Streaming从Kafka生成空Delta记录求助

Kafka流数据摄入Databricks后Delta表仅创建空分区问题

我们在Databricks中基于Kafka实现流数据摄入,运行Spark作业后能看到Parquet分区创建,但分区仅包含元数据,没有实际业务数据。以下是作业代码,求解决方法:

# Databricks notebook source
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
from pyspark.sql.functions import current_timestamp
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, BooleanType, DoubleType, LongType

# COMMAND ----------

spark_depenedencies_jars = [ 
    "org.apache.hadoop:hadoop-common:3.3.1",
    "org.apache.hadoop:hadoop-client:3.3.1",
    "org.apache.hadoop:hadoop-aws:3.3.1",
    "org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.1",
    "org.apache.spark:spark-streaming-kafka-0-10_2.12:3.3.1",
    "io.delta:delta-core_2.12:2.2.0"
]
spark_depenedencies_jars_str = ",".join(spark_depenedencies_jars)

# COMMAND ----------

KAFKA_BOOTSTRAP_SERVERS = "10.0.0.91:9094"
BASE_CHECKPOINT_LOCATION = "wasbs://warehouse@miniobucketsphera.blob.core.windows.net/checkpoint/"
BASE_DELTA_DIR_VM_AZURE = "wasbs://warehouse@miniobucketsphera.blob.core.windows.net/delta/inventory/"

spark = SparkSession.builder \
    .appName("Kafka Streaming Example") \
    .master(SPARK_MASTER_LOCAL) \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
    .config("spark.databricks.delta.retentionDurationCheck.enabled", "false") \
    .config("spark.hadoop.hive.metastore.uris", HIVE_METASTORE_URI) \
    .config("spark.sql.catalogImplementation", "hive") \
    .enableHiveSupport() \
    .getOrCreate()
   
df.write.format("delta").mode("overwrite").save(delta_wasbs_path)

# COMMAND ----------

# .writeStream \ test

def read_kafka_stream(kafka_topic, schema):
    kafka_stream = spark \
                .readStream \
                .format("kafka") \
                .option("kafka.bootstrap.servers", KAFKA_BOOTSTRAP_SERVERS) \
                .option("subscribe", kafka_topic) \
                .option("failOnDataLoss","false") \
                .option("startingOffsets", "earliest") \
                .load()
    data_stream = kafka_stream.selectExpr("cast (value as string) as json") \
                            .select(from_json("json", schema).alias("cdc")) \
                            .select("cdc.payload.after.*", "cdc.payload.op")
    data_stream = data_stream.withColumn("curr_timestamp", current_timestamp())
    return data_stream

# COMMAND ----------

# .save() test

def write_delta_table(data_stream, checkpoint_location, delta_dir):
    data_stream.writeStream \
                .format("delta") \
                .outputMode("append") \
                .option("checkpointLocation", checkpoint_location) \
                .start(delta_dir)

# COMMAND ----------

from pyspark.sql.types import StructType, StructField, StringType, IntegerType

schema_table_schema = StructType([
    StructField("id", IntegerType(), True),
    StructField("age", IntegerType(), True),
    StructField("name", StringType(), True),
])

source_schema = StructType([
    StructField("version", StringType(), False),
    StructField("connector", StringType(), False),
    StructField("name", StringType(), False),
    StructField("ts_ms", IntegerType(), False),
    StructField("snapshot", StringType(), True),
    StructField("db", StringType(), False),
    StructField("sequence", StringType(), True),
    StructField("table", StringType(), True),
    StructField("server_id", IntegerType(), False),
    StructField("gtid", StringType(), True),
    StructField("file", StringType(), False),
    StructField("pos", IntegerType(), False),
    StructField("row", IntegerType(), False),
    StructField("thread", IntegerType(), True),
    StructField("query", StringType(), True),
])

transaction_schema = StructType([
    StructField("id", StringType(), False),
    StructField("total_order", IntegerType(), False),
    StructField("data_collection_order", IntegerType(), False),
])

schema_payload_schema = StructType([
    StructField("before", schema_table_schema, True),
    StructField("after", schema_table_schema, True),
    StructField("source", source_schema, True),
    StructField("op", StringType(), False),
    StructField("ts_ms", IntegerType(), True),
    StructField("transaction", transaction_schema, True),
])

schema_schema = StructType([
    StructField("schema", StringType(), True),
    StructField("payload", schema_payload_schema, True),
])


# COMMAND ----------

schema_table_topic = "sqlf-eus-sccspoc-dev.schema1.table1"


# COMMAND ----------

schema_table_checkpoint_location = BASE_CHECKPOINT_LOCATION + "schema_table"


# COMMAND ----------

schema_table_delta_dir = BASE_DELTA_DIR_VM_AZURE + "schema_table"


# COMMAND ----------

schema_table_stream = read_kafka_stream(schema_table_topic, schema_schema)


# COMMAND ----------

write_delta_table(schema_table_stream, schema_table_checkpoint_location, schema_table_delta_dir)

排查与解决思路

  • 验证Kafka主题数据:先确认目标主题sqlf-eus-sccspoc-dev.schema1.table1存在有效数据,用Kafka命令行工具检查:

    kafka-console-consumer.sh --bootstrap-server 10.0.0.91:9094 --topic sqlf-eus-sccspoc-dev.schema1.table1 --from-beginning
    

    如果没有数据,问题出在数据源端,需排查CDC连接器或数据生产者。

  • 检查Schema匹配度:Kafka消息的JSON结构是否和代码中定义的schema_schema完全一致?临时修改read_kafka_stream函数,打印原始消息:

    def read_kafka_stream(kafka_topic, schema):
        kafka_stream = spark \
                    .readStream \
                    .format("kafka") \
                    .option("kafka.bootstrap.servers", KAFKA_BOOTSTRAP_SERVERS) \
                    .option("subscribe", kafka_topic) \
                    .option("failOnDataLoss","false") \
                    .option("startingOffsets", "earliest") \
                    .load()
        # 打印原始消息到控制台
        kafka_stream.selectExpr("cast(value as string)").writeStream \
            .format("console") \
            .outputMode("append") \
            .start()
        data_stream = kafka_stream.selectExpr("cast (value as string) as json") \
                                .select(from_json("json", schema).alias("cdc")) \
                                .select("cdc.payload.after.*", "cdc.payload.op")
        data_stream = data_stream.withColumn("curr_timestamp", current_timestamp())
        return data_stream
    

    若from_json解析失败会返回null,导致后续无数据输出,需调整Schema匹配实际消息结构。

  • 移除不必要的Spark配置:在Databricks环境中,不需要手动指定master(SPARK_MASTER_LOCAL),移除该配置项,避免本地模式引发的资源或连接问题。

  • 过滤无效数据:CDC消息中op为删除(如"d")时,after字段为null,这类记录会被无意义写入,可添加过滤:

    data_stream = data_stream.filter(col("cdc.payload.after").isNotNull())
    
  • 验证存储权限:确认Databricks集群对Azure Blob存储路径有读写权限,手动写入测试数据验证:

    test_df = spark.createDataFrame([(1, 25, "test", "i", current_timestamp())], ["id", "age", "name", "op", "curr_timestamp"])
    test_df.write.format("delta").mode("append").save(schema_table_delta_dir)
    

    若写入失败,需调整存储账户权限或路径配置。

  • 查看流作业日志:在Databricks作业监控页面查看流作业日志,排查是否有Kafka连接失败、Schema解析错误、存储写入异常等报错信息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 01:02:49