【问题标题】:Structured Streaming exception when using append output mode with watermark使用带有水印的附加输出模式时的结构化流式处理异常
【发布时间】:2018-01-16 14:13:14
【问题描述】:

尽管我使用的是withWatermark(),但我在运行 Spark 作业时收到以下错误消息:

线程“main”中的异常 org.apache.spark.sql.AnalysisException:当流式 DataFrames/DataSets 上存在流式聚合时,不支持附加输出模式而没有水印;;

从我在programming guide 中看到的内容来看,这完全符合预期用途(以及示例代码)。有谁知道可能出了什么问题?

提前致谢!

相关代码(Java 8、Spark 2.2.0):

StructType logSchema = new StructType()
        .add("timestamp", TimestampType)
        .add("key", IntegerType)
        .add("val", IntegerType);

Dataset<Row> kafka = spark
        .readStream()
        .format("kafka")
        .option("kafka.bootstrap.servers", brokers)
        .option("subscribe", topics)
        .load();

Dataset<Row> parsed = kafka
        .select(from_json(col("value").cast("string"), logSchema).alias("parsed_value"))
        .select("parsed_value.*");

Dataset<Row> tenSecondCounts = parsed
        .withWatermark("timestamp", "10 minutes")
        .groupBy(
            parsed.col("key"),
            window(parsed.col("timestamp"), "1 day"))
        .count();

StreamingQuery query = tenSecondCounts
        .writeStream()
        .trigger(Trigger.ProcessingTime("10 seconds"))
        .outputMode("append")
        .format("console")
        .option("truncate", false)
        .start();

【问题讨论】:

    标签: java apache-spark spark-structured-streaming


    【解决方案1】:

    问题出在parsed.col。将其替换为 col 将解决此问题。我建议始终使用col 函数而不是Dataset.col

    Dataset.col 返回resolved columncol 返回unresolved column

    parsed.withWatermark("timestamp", "10 minutes") 将创建一个具有相同名称的新列的新数据集。水印信息附加在新数据集中的timestamp列,而不是parsed.col("timestamp"),因此groupBy中的列没有水印。

    当您使用未解析的列时,Spark 会为您找出正确的列。

    【讨论】:

    猜你喜欢
    • 2019-06-04
    • 2018-07-22
    • 2021-05-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-09-07
    • 2023-03-26
    相关资源
    最近更新 更多