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

