Spark Streaming无法向Iceberg表写入数据排查求助
Spark Streaming写入Iceberg仅生成元数据无数据文件问题
我在本地用Docker部署了Spark-Iceberg、Rest Catalog与MinIO环境。运行Spark Streaming代码时,程序能读取Kafka主题的数据,尝试写入Iceberg表,但MinIO上只生成元数据文件,没有数据文件,而且没报任何错误。已经试过调整输出模式、触发器等参数,还是解决不了。相关代码和Spark提交命令如下:
Python代码
from pyspark.sql import SparkSession from pyspark.sql.functions import col, from_json, min, max, avg, window from pyspark.sql.types import StructField, StructType, DoubleType, StringType, TimestampType if __name__ == "__main__": spark = SparkSession.builder.appName("Kafka-Iceberg-Stream-Processor").getOrCreate() spark.sparkContext.setLogLevel("ERROR") schema = StructType([ StructField("BTC_EUR", DoubleType(), True), StructField("BTC_INR", DoubleType(), True), StructField("BTC_USD", DoubleType(), True), StructField("ETH_EUR", DoubleType(), True), StructField("ETH_INR", DoubleType(), True), StructField("ETH_USD", DoubleType(), True), StructField("timestamp", TimestampType(), True) ]) kafka_df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "192.168.1.16:9092") \ .option("subscribe", "crypto") \ .option("startingOffsets", "earliest") \ .load() parsed_df = kafka_df.selectExpr( "CAST(value AS STRING) AS value", "timestamp" ).withColumn("parsed_json", from_json(col("value"), schema)) \ .select( col("parsed_json.BTC_EUR").alias("BTC_EUR"), col("parsed_json.BTC_INR").alias("BTC_INR"), col("parsed_json.BTC_USD").alias("BTC_USD"), col("parsed_json.ETH_EUR").alias("ETH_EUR"), col("parsed_json.ETH_INR").alias("ETH_INR"), col("parsed_json.ETH_USD").alias("ETH_USD"), col("parsed_json.timestamp").alias("event_time"), col("timestamp").alias("processing_time") ) windowed_df = parsed_df \ .withWatermark("event_time", "1 minute") \ .groupBy(window(col("event_time"), "5 minutes", "5 minutes")) \ .agg( max("BTC_EUR").alias("max_BTC_EUR"), min("BTC_EUR").alias("min_BTC_EUR"), avg("BTC_EUR").alias("avg_BTC_EUR"), max("BTC_INR").alias("max_BTC_INR"), min("BTC_INR").alias("min_BTC_INR"), avg("BTC_INR").alias("avg_BTC_INR"), max("BTC_USD").alias("max_BTC_USD"), min("BTC_USD").alias("min_BTC_USD"), avg("BTC_USD").alias("avg_BTC_USD"), max("ETH_EUR").alias("max_ETH_EUR"), min("ETH_EUR").alias("min_ETH_EUR"), avg("ETH_EUR").alias("avg_ETH_EUR"), max("ETH_INR").alias("max_ETH_INR"), min("ETH_INR").alias("min_ETH_INR"), avg("ETH_INR").alias("avg_ETH_INR"), max("ETH_USD").alias("max_ETH_USD"), min("ETH_USD").alias("min_ETH_USD"), avg("ETH_USD").alias("avg_ETH_USD") ).withColumn("window_start", col("window.start")) \ .withColumn("window_end", col("window.end")) \ .drop("window") \ .select("window_start", "window_end", "max_BTC_EUR", "min_BTC_EUR", "avg_BTC_EUR", "max_BTC_INR", "min_BTC_INR", "avg_BTC_INR", "max_BTC_USD", "min_BTC_USD", "avg_BTC_USD", "max_ETH_EUR", "min_ETH_EUR", "avg_ETH_EUR", "max_ETH_INR", "min_ETH_INR", "avg_ETH_INR", "max_ETH_USD", "min_ETH_USD", "avg_ETH_USD") # Create the table if it doesn't exist spark.sql(""" CREATE TABLE IF NOT EXISTS db.crypto_metrics5 ( window_start TIMESTAMP, window_end TIMESTAMP, max_BTC_EUR DOUBLE, min_BTC_EUR DOUBLE, avg_BTC_EUR DOUBLE, max_BTC_INR DOUBLE, min_BTC_INR DOUBLE, avg_BTC_INR DOUBLE, max_BTC_USD DOUBLE, min_BTC_USD DOUBLE, avg_BTC_USD DOUBLE, max_ETH_EUR DOUBLE, min_ETH_EUR DOUBLE, avg_ETH_EUR DOUBLE, max_ETH_INR DOUBLE, min_ETH_INR DOUBLE, avg_ETH_INR DOUBLE, max_ETH_USD DOUBLE, min_ETH_USD DOUBLE, avg_ETH_USD DOUBLE ) USING iceberg """) print(spark.sql("SHOW TABLES IN rest.db").show()) print(windowed_df.printSchema()) # Write stream data to Iceberg query = windowed_df.writeStream \ .outputMode("complete") \ .format("iceberg") \ .option("checkpointLocation", "/tmp/spark/checkpoints/crypto_metrics5") \ .toTable("rest.db.crypto_metrics5") query.awaitTermination()
Spark提交命令
spark-submit \ --packages org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.7.1,org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.4 \ --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \ --conf spark.sql.catalog.rest=org.apache.iceberg.spark.SparkCatalog \ --conf spark.sql.catalog.rest.type=rest \ --conf spark.sql.catalog.rest.uri=http://rest:8181 \ --conf spark.sql.catalog.rest=org.apache.iceberg.spark.SparkCatalog \ --conf spark.sql.catalog.rest.warehouse=s3://warehouse \ --conf spark.sql.defaultCatalog=rest \ test.py
问题排查与解决方案
1. 窗口计算未触发输出
代码使用5分钟滚动窗口+1分钟水位线,只有当窗口结束时间超过水位线后,Spark才会输出该窗口的聚合结果。如果Kafka数据的event_time未达到窗口结束时间+水位线,就不会生成数据文件。
- 临时将窗口大小改为10秒、水位线改为5秒,快速验证是否有输出;
- 检查Kafka消息中的
timestamp字段是否符合预期时间范围。
2. 补全MinIO认证配置
提交命令重复配置了spark.sql.catalog.rest=org.apache.iceberg.spark.SparkCatalog,且缺失MinIO的S3认证参数,导致Spark无法正常写入数据到MinIO。修正后的提交命令如下(替换实际密钥):
spark-submit \ --packages org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.7.1,org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.4 \ --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \ --conf spark.sql.catalog.rest=org.apache.iceberg.spark.SparkCatalog \ --conf spark.sql.catalog.rest.type=rest \ --conf spark.sql.catalog.rest.uri=http://rest:8181 \ --conf spark.sql.catalog.rest.warehouse=s3://warehouse \ --conf spark.sql.defaultCatalog=rest \ --conf spark.hadoop.fs.s3a.endpoint=http://minio:9000 \ --conf spark.hadoop.fs.s3a.access.key=你的MinIO访问密钥 \ --conf spark.hadoop.fs.s3a.secret.key=你的MinIO秘密密钥 \ --conf spark.hadoop.fs.s3a.path.style.access=true \ --conf spark.hadoop.fs.s3a.impl=org.apache.hadoop.fs.s3a.S3AFileSystem \ test.py
3. 添加触发器强制输出
Iceberg流式写入可通过触发器强制触发输出,修改写入代码:
from pyspark.sql.streaming import Trigger query = windowed_df.writeStream \ .outputMode("complete") \ .format("iceberg") \ .option("checkpointLocation", "/tmp/spark/checkpoints/crypto_metrics5") \ .trigger(Trigger.ProcessingTime("1 minute")) \ .toTable("rest.db.crypto_metrics5")
4. 降低日志级别查看细节
当前日志级别为ERROR,屏蔽了关键提示信息,临时改为INFO级别:
spark.sparkContext.setLogLevel("INFO")
5. 验证Iceberg表存储路径
通过Spark SQL查询表元数据,确认存储路径指向MinIO:
DESCRIBE EXTENDED rest.db.crypto_metrics5;
查看Location字段是否为s3://warehouse/db/crypto_metrics5,并确保MinIO该路径有读写权限。
内容的提问来源于stack exchange,提问作者Rohit Anil
相关产品推荐
相关产品推荐

