PySpark执行顺序因优化异常导致Cassandra更新结果不符合预期问题咨询
看起来你遇到了PySpark惰性执行特性和外部数据源读写交互时的典型坑,我来帮你一步步排查和解决这个问题。
首先,你的核心逻辑是没问题的:读取CSV数据、拉取Cassandra对应日期区间的存量数据、合并后聚合去重、最后写回覆盖Cassandra的旧数据。但结果只有CSV的数据,说明要么存量数据没被正确读取到,要么合并/聚合环节出了问题。
最可能的几个问题点及排查方案
1. 存量数据读取异常(telco_day_aggr为空)
你怀疑写操作先执行覆盖了数据,但从代码顺序看,write_on_cassandra_dev是最后一步,理论上不会提前执行。更可能的是你的日期过滤条件没匹配到Cassandra里的数据:
- 先打印下你拿到的日期范围,确认是否符合预期:
print(f"CSV数据的日期范围:{lower_date} 到 {upper_date}") - 检查Cassandra表中
data字段的类型和df中data的类型是否一致:比如df里是带时分秒的timestamp,Cassandra里是只存日期的date,这时候between条件会因为时间部分不匹配过滤掉所有数据。解决方法是统一类型:# 把CSV的日期转成date类型再取最大最小 max_min_dates = df.agg(F.max(df['data'].cast('date')), F.min(df['data'].cast('date'))).collect()[0] lower_date = max_min_dates[1] upper_date = max_min_dates[0] - 读取Cassandra数据后,强制打印样本和数据量,确认是否真的读到了存量数据:
telco_day_aggr = read_from_cassandra_dev(f'telco_{aggr_map_dict[aggr]}_aggr').where(F.col('data').between(lower_date,upper_date)) print(f"Cassandra读取到的记录数:{telco_day_aggr.count()}") print("Cassandra数据样本:") telco_day_aggr.show(5) print("Cassandra数据结构:") telco_day_aggr.printSchema()
2. Union操作因列匹配失败丢失数据
PySpark的union是按列位置匹配的,而不是列名。如果df和telco_day_aggr的列顺序不一样,会导致数据错位,甚至聚合后看起来只有CSV的数据。解决这个问题的最佳方式是改用unionByName:
# 按列名合并,不允许缺失列,这样列名不匹配会直接报错,方便排查 union_df = df.unionByName(telco_day_aggr, allowMissingColumns=False)
同时还要确认两个DataFrame的列类型完全一致,比如presenze在df里是int,在Cassandra里是long,虽然能合并,但聚合时可能出现异常。
3. 空DataFrame的schema不匹配
当telco_day_aggr为空时,你用create_empty_df()生成空表,但如果这个空表的schema和df不一致,Union后会丢失列或者导致数据异常。确保create_empty_df()返回的DataFrame和df的schema完全相同,可以这样生成:
def create_empty_df(): # 直接基于df的schema生成空表,避免手动定义出错 return spark.createDataFrame([], df.schema)
4. 惰性执行的潜在影响
虽然理论上DAG的依赖关系会保证先读再聚合再写,但为了彻底避免读写顺序的潜在问题,可以强制触发Cassandra数据的读取并缓存:
telco_day_aggr = read_from_cassandra_dev(...).where(...) # 缓存数据并强制加载,避免后续重复读取 telco_day_aggr = telco_day_aggr.cache() # 触发执行 telco_day_aggr.count()
验证步骤
在聚合完成后,先打印聚合结果的样本,确认是否包含CSV和Cassandra的合并数据:
print("聚合后最终数据样本:") output_df.show(10)
如果这里能看到合并的数据,那问题就出在写Cassandra的环节;如果看不到,就回到前面的步骤排查读取和合并问题。
备注:内容来源于stack exchange,提问作者Gabriele Sciurti

