如何在PySpark中获取CDF Merge操作影响的行数?
PySpark中获取Delta Merge操作的影响行数
实时获取Merge执行统计
Delta Lake的merge().execute()方法会返回包含操作统计的结果对象,直接捕获并提取即可得到插入、更新(若包含删除逻辑则同步获取删除数)的行数:
from delta.tables import * deltaTablePeople = DeltaTable.forPath(spark, '/tmp/delta/people-10m') deltaTablePeopleUpdates = DeltaTable.forPath(spark, '/tmp/delta/people-10m-updates') dfUpdates = deltaTablePeopleUpdates.toDF() # 执行Merge并捕获结果 merge_outcome = deltaTablePeople.alias('people') \ .merge( dfUpdates.alias('updates'), 'people.id = updates.id' ) \ .whenMatchedUpdate(set = { "id": "updates.id", "firstName": "updates.firstName", "middleName": "updates.middleName", "lastName": "updates.lastName", "gender": "updates.gender", "birthDate": "updates.birthDate", "ssn": "updates.ssn", "salary": "updates.salary" }) \ .whenNotMatchedInsert(values = { "id": "updates.id", "firstName": "updates.firstName", "middleName": "updates.middleName", "lastName": "updates.lastName", "gender": "updates.gender", "birthDate": "updates.birthDate", "ssn": "updates.ssn", "salary": "updates.salary" }) \ .execute() # 解析统计数据 stats = merge_outcome.asDict() updated_rows = stats['numUpdatedRows'] inserted_rows = stats['numInsertedRows'] # 若包含whenMatchedDelete逻辑,可通过stats['numDeletedRows']获取删除行数 print(f"更新行数: {updated_rows}, 插入行数: {inserted_rows}")
事后查询历史操作统计
如果需要追溯过往Merge操作的影响行数,可通过Delta表的版本历史记录查询:
# 获取表的版本历史 history_df = deltaTablePeople.history() # 筛选最新的Merge操作记录 latest_merge_record = history_df.filter(history_df['operation'] == 'MERGE') \ .orderBy(history_df['version'].desc()) \ .first() # 提取操作指标 operation_metrics = latest_merge_record['operationMetrics'] print(f"更新行数: {operation_metrics['numUpdatedRows']}, 插入行数: {operation_metrics['numInsertedRows']}")
两种方法按需选择:前者适合Merge执行后实时获取统计数据,后者用于事后核查历史操作的影响情况。
内容的提问来源于stack exchange,提问作者Aditya Raj
相关产品推荐
相关产品推荐

