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

PySpark无水印流聚合写入Cassandra遇输出模式兼容问题

问题描述

数据与需求

Kafka Topic流式传输用户数据,每条消息格式如下:

message1 = {'username': 'user1', 'account_balance': 25.0}
message2 = {'username': 'user2', 'account_balance': 2.0}
message3 = {'username': 'user2', 'account_balance': 4.0}
message4 = {'username': 'user1', 'account_balance': 27.0}

使用PySpark结构化流读取数据,需要计算截至当前的account_balance平均值,预期每接收一条消息后的聚合结果:

agg_df1 = {'username': 'user1', 'account_balance': 25.0}
agg_df2 = {'username': 'user1', 'account_balance': 25.0}, {'username': 'user2', 'account_balance': 2.0}
agg_df3 = {'username': 'user1', 'account_balance': 25.0}, {'username': 'user2', 'account_balance': 3.0}
agg_df4 = {'username': 'user1', 'account_balance': 26.0}, {'username': 'user2', 'account_balance': 2.0}

最终需要将每个聚合结果(带时间戳)附加写入Cassandra表。

尝试的代码与错误

第一次尝试(Append模式)

编写的PySpark脚本:

# parsed_df 是解析Kafka消息后的DataFrame
# parsed_df = DataFrame[username: string, account_balance: float]

summary_df = parsed_df.groupBy("username") \
.agg(avg("account_balance").alias("account_balance_avg")) \
.withColumn("uuid", make_uuid()) \
.withColumn("ingest_timestamp", current_timestamp())

summary_df.writeStream.outputMode("append").format("org.apache.spark.sql.cassandra") \
  .option("checkpointLocation", "/tmp/check_point/") \
  .options(table="avg_account_balance", keyspace="user") \
  .start()

抛出错误:

Traceback (most recent call last):
  File "<stdin>", line 4, in <module>
  File "/opt/spark/python/pyspark/sql/streaming/readwriter.py", line 1385, in start
    return self._sq(self._jwrite.start())
  File "/opt/spark/python/lib/py4j-0.10.9.7-src.zip/py4j/java_gateway.py", line 1322, in __call__
  File "/opt/spark/python/pyspark/errors/exceptions/captured.py", line 175, in deco
    raise converted from None
pyspark.errors.exceptions.captured.AnalysisException: Append output mode not supported when there are streaming aggregations on streaming DataFrames/DataSets without watermark;

核心矛盾:不想使用watermark,因为需要计算从开始到现在的全量平均值。

第二次尝试(Complete模式)

改用outputMode("complete")写入Cassandra,抛出错误:

java.lang.UnsupportedOperationException: You are attempting to use overwrite mode which will truncate
this table prior to inserting data. If you would merely like
to change data already in the table use the "Append" mode.
To actually truncate please pass in true value to the option
"confirm.truncate" or set that value to true in the session.conf
 when saving.

但用complete模式写入控制台是正常的,代码如下:

parsed_df.groupBy("username") \
.agg(avg("account_balance").alias("account_balance_avg")) \
.withColumn("uuid", make_uuid()) \
.withColumn("ingest_timestamp", current_timestamp()) \
.writeStream.outputMode("complete").format("console").start()
解决方案

思路

要保留全量历史的聚合结果(每次更新后把完整聚合结果追加到Cassandra,而非覆盖),但complete模式会触发Cassandra的覆盖行为,append模式在无watermark的流式聚合中不被支持。可以通过foreachBatch自定义写入逻辑,将complete模式的输出转换为追加写入的流。

代码实现

from pyspark.sql.functions import avg, current_timestamp, expr

def write_to_cassandra(batch_df, batch_id):
    # 每个微批的聚合结果以append模式写入Cassandra
    batch_df.write.format("org.apache.spark.sql.cassandra") \
        .options(table="avg_account_balance", keyspace="user") \
        .mode("append") \
        .save()

# 构建聚合逻辑
summary_df = parsed_df.groupBy("username") \
    .agg(avg("account_balance").alias("account_balance_avg")) \
    .withColumn("uuid", expr("uuid()")) \
    .withColumn("ingest_timestamp", current_timestamp())

# 使用foreachBatch自定义写入流程
summary_df.writeStream \
    .outputMode("complete") \
    .foreachBatch(write_to_cassandra) \
    .option("checkpointLocation", "/tmp/check_point/") \
    .start() \
    .awaitTermination()

说明

  • foreachBatch允许在每个微批处理时自定义写入逻辑,这里将complete模式输出的全量聚合结果以append模式写入Cassandra,既保留了全量平均值的计算,又实现了追加写入的需求。
  • 需确保Cassandra表的主键包含uuid(每次聚合都会生成新的uuid),避免主键冲突。示例表结构:
CREATE TABLE user.avg_account_balance (
    uuid UUID PRIMARY KEY,
    username TEXT,
    account_balance_avg DOUBLE,
    ingest_timestamp TIMESTAMP
);

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 20:25:07