多会话从S3导入Redshift出现重复行与缺失行问题排查
多会话S3到Redshift导入出现重复/行缺失问题
我在执行多会话从S3向Redshift导入数据时,遇到了数据重复和行缺失的问题。编写了启动5个顺序执行会话的测试脚本后,频繁复现了重复问题。我怀疑自己的COPY/IMPORT实现存在问题,以下是核心流程代码:
LOGGER.debug('Importing into %s', tmp_table) cursor.execute('BEGIN;') cursor.execute('LOCK {schema}.{table} ;'.format(**sql_args)) cursor.execute('CREATE TEMP TABLE {tmp_table} ( LIKE {schema}.{table} );'.format(**sql_args)) LOGGER.debug('{s3file}'.format(**sql_args)) cursor.execute('COPY {tmp_table} FROM \'s3://{bucket}/{s3file}\' CREDENTIALS \'aws_access_key_id={s3user};aws_secret_access_key={s3pwd}\' {copy_options};'.format(**sql_args)) cursor.execute('SELECT COUNT(*) FROM {tmp_table}'.format(**sql_args)) res = cursor.fetchall() cursor.execute('INSERT INTO {schema}.{table} SELECT * FROM {tmp_table};'.format(**sql_args)) LOGGER.debug('Dropping {tmp_table}'.format(**sql_args)) cursor.execute('DROP TABLE {tmp_table};'.format(**sql_args)) cursor.execute('END;'.format(**sql_args))
我已尝试调整锁机制和自动提交设置,但问题仍未解决,恳请帮助。
问题诊断与修复方案
1. 事务提交错误(核心问题)
Redshift不支持使用END;提交事务,正确的事务结束命令是COMMIT;。你的代码中用END;无法完成事务提交,会导致两种情况:
- 若客户端自动提交关闭,会话结束时Redshift会自动回滚未提交的事务,引发行缺失;
- 若自动提交设置异常,可能导致部分操作被重复执行,引发数据重复。
修复:将cursor.execute('END;')替换为cursor.execute('COMMIT;'),同时添加异常捕获逻辑,在出错时执行ROLLBACK;。
2. 临时表未继承唯一约束
CREATE TEMP TABLE ... LIKE ...默认不会复制原表的主键、唯一约束,导致临时表中可能存在重复行。如果目标表也没有唯一约束,插入后就会产生重复数据。
修复:修改创建临时表的语句,添加INCLUDING CONSTRAINTS参数,让临时表继承原表的约束:
CREATE TEMP TABLE {tmp_table} ( LIKE {schema}.{table} INCLUDING CONSTRAINTS );
3. 未处理COPY异常与数据验证
你的代码没有捕获COPY命令的异常,也没有验证临时表的导入行数是否符合预期。如果COPY失败(比如数据格式错误),空的临时表会被插入到目标表,引发行缺失;若重试机制触发,可能重复导入同一份文件,引发数据重复。
修复:添加异常处理和数据验证逻辑,示例如下:
try: cursor.execute('BEGIN;') cursor.execute('LOCK {schema}.{table} ;'.format(**sql_args)) # 创建带约束的临时表 cursor.execute('CREATE TEMP TABLE {tmp_table} ( LIKE {schema}.{table} INCLUDING CONSTRAINTS );'.format(**sql_args)) LOGGER.debug('Processing s3://{bucket}/{s3file}'.format(**sql_args)) cursor.execute('COPY {tmp_table} FROM \'s3://{bucket}/{s3file}\' CREDENTIALS \'aws_access_key_id={s3user};aws_secret_access_key={s3pwd}\' {copy_options};'.format(**sql_args)) # 验证临时表导入行数 cursor.execute('SELECT COUNT(*) FROM {tmp_table}'.format(**sql_args)) tmp_count = cursor.fetchone()[0] if tmp_count == 0: LOGGER.warning('No data copied from s3://{bucket}/{s3file}'.format(**sql_args)) cursor.execute('ROLLBACK;') return # 插入目标表并验证行数 cursor.execute('INSERT INTO {schema}.{table} SELECT * FROM {tmp_table};'.format(**sql_args)) cursor.execute('GET DIAGNOSTICS insert_count = ROW_COUNT;') insert_count = cursor.fetchone()[0] if insert_count != tmp_count: LOGGER.error('Insert mismatch: copied %d rows, inserted %d rows', tmp_count, insert_count) cursor.execute('ROLLBACK;') return cursor.execute('DROP TABLE {tmp_table};'.format(**sql_args)) cursor.execute('COMMIT;') except Exception as e: LOGGER.error('Import failed: %s', str(e)) cursor.execute('ROLLBACK;') raise
4. 其他建议
- 确保测试脚本中每个会话处理的S3文件唯一,避免同一份文件被重复导入;
- 检查目标表是否有主键或唯一约束,从源头上阻止重复数据插入;
- 确认Redshift的锁机制生效:
LOCK {schema}.{table}默认是ACCESS EXCLUSIVE锁,顺序执行的会话会排队获取锁,保证INSERT操作串行执行。
内容的提问来源于stack exchange,提问作者jimhens
相关产品推荐
相关产品推荐

