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

如何实现从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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 21:13:13