如何实现从SQL表循环批量读取10000条事件且不丢失数据
实现ClickHouse循环分批读取日志数据(每次10000条,无数据丢失)
核心思路
- 记录上次读取的最后时间点
last_point,每次从该时间点之后读取数据 - 每次查询按
My_time升序排序,保证读取顺序和数据写入顺序一致 - 循环执行读取逻辑,某次读取不到新数据时,等待一段时间后重试,避免空轮询浪费资源
原代码存在的问题
last_value(My_time)用法错误:last_value是窗口函数,必须配合OVER()子句使用,原写法无法正确获取指定范围的最后时间- 缺少循环逻辑:仅执行单次查询,无法实现持续分批读取
- 连接字符串语法错误:
localhost: local_port应改为localhost:{local_port},变量未正确嵌入
修正后的完整代码
import time import pandas as pd from sshtunnel import SSHTunnelForwarder from sqlalchemy import create_engine # 配置参数 ssh_host = '10.x.x.x' ssh_port = 22 ssh_username = 'your_ssh_username' ssh_password = 'your_ssh_password' click_username = 'your_clickhouse_username' click_password = 'your_clickhouse_password' remote_db_port = 9000 db_name = 'My_table' table_name = 'event' batch_size = 10000 wait_interval = 60 # 无新数据时等待60秒再重试 with SSHTunnelForwarder( (ssh_host, ssh_port), ssh_username=ssh_username, ssh_password=ssh_password, remote_bind_address=('localhost', remote_db_port) ) as server: local_port = server.local_bind_port # 修正连接字符串:正确嵌入local_port变量 engine = create_engine(f'clickhouse://{click_username}:{click_password}@localhost:{local_port}/{db_name}') # 初始化上次读取的时间点:先获取表中最早的时间,无数据则设为初始时间 with engine.connect() as conn: init_query = f"SELECT MIN(My_time) as min_time FROM {table_name}" init_result = pd.read_sql(init_query, conn) last_point = init_result['min_time'].iloc[0] if not init_result.empty else '1970-01-01 00:00:00' print(f"开始分批读取数据,初始读取起点:{last_point}") while True: # 参数化查询:从last_point之后读取最多batch_size条,按时间升序 query = f""" SELECT * FROM {table_name} WHERE My_time > %(last_point)s ORDER BY My_time ASC LIMIT {batch_size} """ # 读取数据 batch_data = pd.read_sql(query, engine, params={'last_point': last_point}) if batch_data.empty: print(f"当前无新数据,等待{wait_interval}秒后重试...") time.sleep(wait_interval) continue # 处理当前批次数据(替换成你的业务逻辑) print(f"读取到{len(batch_data)}条数据,时间范围:{batch_data['My_time'].min()} 至 {batch_data['My_time'].max()}") # 示例:打印前5条数据 print(batch_data.head()) # 更新last_point为当前批次的最后一条记录的时间 last_point = batch_data['My_time'].iloc[-1] # 如果本次读取的记录数小于batch_size,说明当前没有更多数据,等待后重试 if len(batch_data) < batch_size: print(f"本次读取记录数不足{batch_size},等待{wait_interval}秒后检查新数据...") time.sleep(wait_interval)
关键注意事项
- 一定要用参数化查询:通过
params传递last_point,既避免SQL注入,又能保证时间类型被正确解析 - 必须按
My_time升序排序:确保每次读取的都是未处理的新数据,不会出现遗漏或重复 - 循环终止逻辑:示例是无限循环,你可以根据需求加终止条件,比如设置最大循环次数、监听外部停止信号等
- 处理重复时间:如果
My_time有重复值,可能会重复读取数据,这时候可以结合表的自增主键一起筛选,比如WHERE (My_time > %(last_point)s) OR (My_time = %(last_point)s AND id > %(last_id)s) - 异常捕获:实际使用时记得加异常处理,捕获数据库连接、查询时的错误,防止程序突然崩溃
内容的提问来源于stack exchange,提问作者Polina
相关产品推荐
相关产品推荐

