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

使用Spark JDBC读取PostgreSQL时分区列致数据重复或缺失

Spark JDBC读取动态PostgreSQL表的数据重复/缺失问题解决

问题场景

用Spark JDBC通过10个连接读取每秒新增数据的PostgreSQL动态表时,先调用calcCountRows统计行数作为upperBound,再用ROW_NUMBER() OVER(ORDER BY NULL)生成的RNO作为分区列做数据读取,结果出现数据重复或缺失,且重复内容每次读取都可能不一样。相关代码如下:

val rows = calcCountRows(table)
val numPartitions = math.min(partitions, rows)

spark.read
  .format("jdbc")
  .options(database.connectionDetails())
  .option("fetchsize", 100)
  .option("dbtable",
    s"""
       |(select ROW_NUMBER() OVER(ORDER BY NULL) AS RNO, subQuery.* from $table subQuery) as "$table"
       |""".stripMargin)
  .option("partitionColumn", "RNO")
  .option("lowerBound", 0)
  .option("upperBound", rows)
  .option("numPartitions", numPartitions)
  .load()

问题根因

  1. ORDER BY NULL导致行号无固定顺序:PostgreSQL里ORDER BY NULL不会保证行的排列顺序,每次查询时行的顺序都可能变化,使得ROW_NUMBER()生成的RNO对应的实际数据行不一致,多个分区可能读到同一条数据,或漏掉某些数据。
  2. 行数统计与读取存在时间差:calcCountRows统计完行数后,表还在持续新增数据,实际读取时的行数已经大于统计值,导致部分新增数据读不到;如果期间有数据删除,还会出现分区范围对应行不存在的缺失情况。
  3. 分区依赖不稳定的临时列:基于动态变化的RNO做分区,每个分区的查询范围对应的实际数据会因为行排序变化而混乱。

解决办法

1. 用表中稳定的唯一有序列做分区键

放弃临时生成的RNO,直接用表自带的唯一且有序列(比如自增主键id、插入时间戳create_time)作为分区列。以自增主键为例:

// 获取主键的最小、最大值
val minId = jdbcTemplate.queryForObject(s"SELECT MIN(id) FROM $table", classOf[Long])
val maxId = jdbcTemplate.queryForObject(s"SELECT MAX(id) FROM $table", classOf[Long])
val numPartitions = math.min(partitions, (maxId - minId + 1).toInt)

spark.read
  .format("jdbc")
  .options(database.connectionDetails())
  .option("fetchsize", 100)
  .option("dbtable", table)
  .option("partitionColumn", "id")
  .option("lowerBound", minId)
  .option("upperBound", maxId)
  .option("numPartitions", numPartitions)
  .load()

如果用时间戳列,要确保该时间戳是记录插入时的固定值,不会被修改,再按时间范围分区。

2. 用事务快照保证数据一致性

如果需要读取某个时间点的完整一致数据,可以利用PostgreSQL的事务快照,让统计行数和读取数据在同一个事务中执行,保证两次操作看到的是同一版本的数据:

// 开启JDBC事务
val conn = DriverManager.getConnection(database.connectionUrl, database.user, database.password)
conn.setAutoCommit(false)
// 保持可重复读隔离级别(PostgreSQL默认)
conn.setTransactionIsolation(Connection.TRANSACTION_REPEATABLE_READ)

// 同一事务内统计行数
val stmt = conn.createStatement()
val rs = stmt.executeQuery(s"SELECT COUNT(*) FROM $table")
rs.next()
val rows = rs.getLong(1)
rs.close()
stmt.close()

// 传递事务连接给Spark
val props = new Properties()
props.putAll(database.connectionDetails())
props.put("connection", conn)

val numPartitions = math.min(partitions, rows.toInt)

val df = spark.read
  .format("jdbc")
  .options(props)
  .option("fetchsize", 100)
  .option("dbtable",
    s"""
       |(select ROW_NUMBER() OVER(ORDER BY id) AS RNO, subQuery.* from $table subQuery) as "$table"
       |""".stripMargin) // 改用主键排序生成稳定RNO
  .option("partitionColumn", "RNO")
  .option("lowerBound", 1)
  .option("upperBound", rows)
  .option("numPartitions", numPartitions)
  .load()

// 读取完成后提交事务
conn.commit()
conn.close()

3. 改用增量读取策略

如果是持续同步动态表,建议每次只读取上次读取之后新增的数据。比如记录上次读取的最大主键或最新时间戳,下次从该值开始读取:

// 从检查点获取上次读取的最大id
val lastMaxId = getLastMaxIdFromCheckpoint()
val currentMaxId = jdbcTemplate.queryForObject(s"SELECT MAX(id) FROM $table", classOf[Long])

val df = spark.read
  .format("jdbc")
  .options(database.connectionDetails())
  .option("fetchsize", 100)
  .option("dbtable", s"SELECT * FROM $table WHERE id > $lastMaxId AND id <= $currentMaxId")
  .option("numPartitions", partitions)
  .load()

// 更新检查点的最大id
updateCheckpoint(currentMaxId)

这种方式从根源上避免了全表读取的一致性问题,适合持续同步场景。

4. 给ROW_NUMBER()指定稳定排序键

如果一定要用ROW_NUMBER()生成分区列,必须指定稳定的排序键(比如主键、唯一索引列),不能用ORDER BY NULL,修改子查询:

(select ROW_NUMBER() OVER(ORDER BY id) AS RNO, subQuery.* from $table subQuery) as "$table"

但这种方式仍需配合事务快照解决统计行数和读取的时间差问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 15:56:03