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
相关产品推荐
相关产品推荐

