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

每日导入循环复用ID的扫描事件数据效率优化问询

咱们先拆解下你的问题,然后一步步给出最有效的提速方案。你现在的核心痛点是Postgres里的写入后DELETE操作效率极低,又没法用临时表或Upsert,再加上scaneventID每45天会复用。下面是几个可落地的解决办法:

方案1:在PySpark层完成全量去重,只写入需要保留的记录

这是最直接的优化思路——把去重逻辑提前到Glue的Spark处理阶段,从根源上减少写入Postgres的数据量,彻底避免后续的DELETE操作。毕竟Spark处理大数据的分布式能力远强于Postgres单表的删除性能。

具体步骤:

  1. 先过滤近45天的数据(因为scaneventID每45天复用,这个窗口内的ID不会和历史周期混淆)
  2. 按scaneventID分组,计算每个ID对应的最小scandatetime(Intake状态)和最大scandatetime(Processed状态)
  3. 筛选出符合这两个时间条件的记录,合并成最终数据集
  4. 直接写入Postgres,全程不需要后续去重

PySpark代码示例:

from pyspark.sql import functions as F

# 读取S3上的Parquet原始数据
raw_df = spark.read.parquet("s3://your-bucket/path/to/scan-events/")

# 过滤近45天的数据(根据scandatetime的实际类型调整,这里假设是timestamp)
cutoff_time = F.current_timestamp() - F.expr("INTERVAL 45 DAYS")
filtered_df = raw_df.filter(F.col("scandatetime") >= cutoff_time)

# 按scaneventID分组,提取每个ID的目标时间点
aggregated_df = filtered_df.groupBy("scaneventID") \
    .agg(
        F.min(F.when(F.col("status") == "Intake", F.col("scandatetime"))).alias("min_intake_time"),
        F.max(F.when(F.col("status") == "Processed", F.col("scandatetime"))).alias("max_processed_time")
    )

# 筛选出需要保留的记录:Intake状态的最小时间、Processed状态的最大时间
final_df = filtered_df.join(aggregated_df, on="scaneventID") \
    .filter(
        (F.col("status") == "Intake") & (F.col("scandatetime") == F.col("min_intake_time")) |
        (F.col("status") == "Processed") & (F.col("scandatetime") == F.col("max_processed_time"))
    ) \
    .drop("min_intake_time", "max_processed_time")

# 优化写入Postgres的参数,提升速度
final_df.write \
    .format("jdbc") \
    .option("url", "jdbc:postgresql://your-postgres-host:5432/your-db") \
    .option("dbtable", "your-target-table") \
    .option("user", "db-user") \
    .option("password", "db-pass") \
    .option("batchsize", "10000")  # 增大批量写入的大小,默认1000太小
    .option("rewriteBatchedStatements", "true")  # 合并批量插入语句,减少网络开销
    .mode("append") \
    .save()

方案2:构造逻辑唯一键,用Postgres的ON CONFLICT实现无Upsert的去重

虽然你说没有唯一标识,但可以结合scaneventID + status + scandatetime构造一个逻辑唯一键——在近45天的窗口内,每个ID的Intake最小时间和Processed最大时间是唯一的,不会重复。

步骤:

  1. 先在Postgres目标表上创建唯一约束:
ALTER TABLE your-target-table 
ADD CONSTRAINT unique_scan_event 
UNIQUE (scaneventID, status, scandatetime);
  1. 在Glue写入时,通过JDBC参数启用ON CONFLICT DO NOTHING,这样重复的记录会被自动跳过,不需要后续DELETE:
final_df.write \
    .format("jdbc") \
    .option("url", "jdbc:postgresql://your-postgres-host:5432/your-db") \
    .option("dbtable", "your-target-table") \
    .option("user", "db-user") \
    .option("password", "db-pass") \
    .option("batchsize", "10000")
    .option("properties", "stringtype=unspecified;")
    .option("createTableOptions", "ON CONFLICT (scaneventID, status, scandatetime) DO NOTHING")
    .mode("append") \
    .save()

如果后续有同一ID的Processed时间更新(比如出现更大的时间),可以把DO NOTHING改成DO UPDATE SET scandatetime = EXCLUDED.scandatetime,确保保留最新的最大时间。

方案3:优化Glue到Postgres的写入性能(如果必须保留现有流程)

如果暂时没法调整数据处理逻辑,至少可以优化写入环节的速度:

  • 增大batchsize参数:从默认的1000调整到10000甚至更高(根据Postgres的max_connections和work_mem配置调整)
  • 启用rewriteBatchedStatements:这个参数会把多条INSERT合并成批量插入,大幅减少网络往返次数
  • 使用分区表:如果Postgres目标表按scandatetime分区,Spark可以直接写入对应分区,减少表级锁的竞争
  • 升级Glue作业资源:增加DPU数量,或者使用G.1X/G.2X等更大的Worker类型,提升Spark的处理速度

方案4:低代码选项——用AWS Glue DataBrew预处理数据

如果不想写大量PySpark代码,DataBrew是个不错的低代码工具:

  1. 创建DataBrew数据集,指向S3上的Parquet扫描事件文件
  2. 构建数据处理配方:
    • 添加过滤步骤:保留近45天的数据
    • 添加分组聚合步骤:按scaneventID分组,提取Intake的最小时间和Processed的最大时间
    • 添加筛选步骤:只保留符合上述时间条件的记录
  3. 配置DataBrew作业,直接把处理后的数据写入Postgres,全程不需要手动写代码

总结

最推荐的是方案1,它在数据写入数据库前就完成了去重,彻底消除了低效的POSTGRES DELETE操作,利用Spark的分布式能力处理数十万条数据非常高效。如果需要更灵活的更新逻辑,可以结合方案2的唯一约束和ON CONFLICT机制。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 08:57:44