首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >RPA+数据库:批量处理百万条数据,比脚本更稳更快

RPA+数据库:批量处理百万条数据,比脚本更稳更快

原创
作者头像
用户12579380
发布2026-07-20 14:30:30
发布2026-07-20 14:30:30
810
举报

一、当百万级数据成为日常

去年双十一,公司电商后台的订单数据表突破了180万条记录。运营同事凌晨三点打电话过来,说需要把过去三个月的订单按省份、品类、支付方式做交叉统计,生成一份给老板的汇报材料。

我第一反应是写个Python脚本跑一宿。结果脚本跑到第47万条的时候,因为某个字段的编码格式不一致,整个进程崩了。重新跑?前面47万条白跑了。不跑?明天上午老板要看。

那天晚上,我盯着报错日志发呆。后来换了个思路——用流程自动化工具搭了一套带异常捕获、断点续跑、分批提交的处理流程。不仅跑完了,还顺手把结果自动导成了Excel发到了钉钉群。

这件事让我重新思考:当数据量从"万级"跃升到"百万级",脚本和自动化工具到底差在哪?


二、百万级数据处理的三大痛点

2.1 连接池撑不住:一次性加载百万条的代价

传统脚本一次性拉取百万条数据,内存直接爆炸。就算用分页查询,频繁的连接建立和断开也会让数据库压力飙升。很多开发者习惯写这样的代码:

制# 典型的问题写法:一次性全量加载 cursor.execute("SELECT * FROM orders WHERE create_time > '2026-04-01'") all_data = cursor.fetchall() # 内存瞬间飙到几个G

这种写法在百万级数据处理场景下就是一颗定时炸弹。更麻烦的是,脚本跑崩了没有任何恢复机制,只能从头再来。

2.2 异常中断无恢复:一条脏数据毁掉整批任务

脚本一旦遇到脏数据、网络抖动、字段类型不匹配,轻则跳过错误记录,重则整个任务失败。大批量数据处理里哪怕只有0.01%的异常,也意味着100条记录可能让整批任务功亏一篑。

而且脚本的异常处理全靠手写try-except,漏一个边界情况就可能全盘皆输。对于高并发数据处理场景来说,这种脆弱性是不可接受的。

2.3 跨系统协同困难:脚本不是万能的

真实业务场景从来不是"查完数据库就结束"。数据要清洗、要格式转换、要对接ERP、要发邮件通知、要生成可视化报表。纯脚本方案需要写大量胶水代码,维护成本极高。

比如你要从MySQL读取数据,清洗后写入PostgreSQL,再生成Excel发邮件——这一套下来,脚本代码可能上千行。而流程自动化软件通过可视化编排,几小时就能搭好,还能支持API触发定时执行,设定好规则自动跑,不用人工守着。


三、RPA+数据库:一种更工程化的自动化数据处理方案

3.1 什么是RPA+数据库方案

简单来说,就是把数据库操作(查询、插入、更新、删除)嵌入到自动化流程中,利用流程引擎的调度能力、容错机制和可视化编排,替代传统的手写脚本。

这种RPA数据库方案的核心优势在于:流程即代码,但比代码更 resilient(有韧性)。对于中小企业自动化个人开发者来说,选择一款合适的RPA软件,往往比从头造轮子更务实。

3.2 架构设计:分层解耦

一个稳健的数据批量处理方案,建议按以下三层设计:

代码语言:javascript
复制
┌─────────────────────────────────────────┐
│  调度层:定时触发 / API触发 / 手动触发    │
├─────────────────────────────────────────┤
│  处理层:分批读取 → 数据清洗 → 批量写入  │
├─────────────────────────────────────────┤
│  持久层:源数据库 → 临时表 → 目标库/文件  │
└─────────────────────────────────────────┘

这里多说一句调度层的设计。好的流程自动化工具不仅支持API触发,还能支持定时执行,比如每天凌晨2点自动跑批。更进阶的还能通过Agent功能,在钉钉、飞书、企微、个人微信里直接发指令触发流程,完成后把结果回调到群里。这种Agent自动化的玩法,让非技术人员也能轻松调度企业数据自动化任务。


四、实战:百万订单数据迁移与清洗

4.1 场景描述

假设有一个电商系统,订单表orders积累了120万条历史数据。需要:

  1. 筛选出2026年Q2的订单(约80万条)
  2. 清洗手机号格式(去除空格、统一+86前缀)
  3. 按省份分组统计销售额
  4. 结果写入新表orders_summary
  5. 生成Excel报表并邮件发送

这是一个典型的数据库数据迁移+数据清洗自动化场景。如果用纯脚本做,开发和调试可能要2-3天。但用RPA+AI的思路,配合可视化编排,2-4小时就能搭好一套稳定的流程。

4.2 传统脚本方案(及其隐患)

代码语言:javascript
复制
import pymysql
import pandas as pd
from datetime import datetime

def migrate_orders():
    conn = pymysql.connect(host='localhost', user='root', 
                           password='xxx', database='shop', charset='utf8mb4')
    try:
        # 隐患1:一次性加载80万条,内存风险
        df = pd.read_sql("""
            SELECT * FROM orders 
            WHERE create_time BETWEEN '2026-04-01' AND '2026-06-30'
        """, conn)

        # 隐患2:无异常隔离,一条脏数据全崩
        df['phone'] = df['phone'].apply(clean_phone)

        # 隐患3:单条插入,性能极差
        for _, row in df.iterrows():
            cursor.execute("INSERT INTO orders_summary ...", row)

    except Exception as e:
        print(f"任务失败: {e}")  # 隐患4:失败即结束,无断点续跑
        conn.rollback()
    finally:
        conn.close()

这段代码的问题很明显:数据库批量插入用的是单条循环,性能极差;没有断点续跑机制,崩了就得重来;更没有脏数据隔离,一条异常记录可能毁掉整批百万条数据批量处理任务。

4.3 流程自动化方案:分批+容错+自愈

下面是一个经过工程化改造的完整数据库批量处理方案。核心思路:分批读取、批量提交、异常隔离、断点续跑

第一步:数据库连接配置

代码语言:javascript
复制
# db_config.py
DB_CONFIG = {
    'host': 'localhost',
    'port': 3306,
    'user': 'rpa_user',
    'password': 'your_password',
    'database': 'shop',
    'charset': 'utf8mb4',
    'connect_timeout': 30,
    'read_timeout': 300,
    'write_timeout': 300
}

# 连接池配置,避免频繁创建连接
POOL_CONFIG = {
    'pool_name': 'rpa_pool',
    'pool_size': 5,
    'pool_reset_session': True
}
第二步:分批读取(游标方式)

代码语言:javascript
复制
import mysql.connector
from mysql.connector import pooling

def create_connection_pool():
    """创建数据库连接池,复用连接降低开销"""
    return mysql.connector.pooling.MySQLConnectionPool(
        pool_name="batch_pool",
        pool_size=5,
        **DB_CONFIG
    )

def fetch_in_batches(pool, batch_size=5000):
    """
    使用服务器端游标分批读取,避免客户端内存爆炸
    每次只加载 batch_size 条记录到内存
    """
    conn = pool.get_connection()
    # SSCursor (Server Side Cursor) 关键:数据不一次性拉到客户端
    cursor = conn.cursor(dictionary=True, buffered=False)

    try:
        cursor.execute("""
            SELECT id, order_no, phone, province, amount, create_time 
            FROM orders 
            WHERE create_time BETWEEN '2026-04-01' AND '2026-06-30'
            ORDER BY id
        """)

        batch = []
        for row in cursor:
            batch.append(row)
            if len(batch) >= batch_size:
                yield batch
                batch = []

        if batch:  # 最后一批
            yield batch

    finally:
        cursor.close()
        conn.close()
第三步:数据清洗(带异常隔离)

代码语言:javascript
复制
import re
import logging

logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s - %(levelname)s - %(message)s',
    handlers=[
        logging.FileHandler('data_clean.log', encoding='utf-8'),
        logging.StreamHandler()
    ]
)

def clean_phone(phone):
    """清洗手机号:去空格、统一格式"""
    if not phone:
        return None
    # 去除所有非数字字符
    digits = re.sub(r'\D', '', str(phone))
    # 统一添加+86前缀(如果没有)
    if len(digits) == 11 and digits.startswith('1'):
        return f'+86{digits}'
    elif len(digits) == 13 and digits.startswith('86'):
        return f'+{digits}'
    return digits  # 异常格式保留原样,后续人工复核

def process_batch(batch_records):
    """
    处理一批数据,单条异常不影响整批
    返回: (成功列表, 失败列表)
    """
    success = []
    failed = []

    for record in batch_records:
        try:
            cleaned = {
                'id': record['id'],
                'order_no': record['order_no'],
                'phone': clean_phone(record['phone']),
                'province': record['province'] or '未知',
                'amount': float(record['amount']) if record['amount'] else 0.0,
                'create_time': record['create_time']
            }
            success.append(cleaned)
        except Exception as e:
            failed.append({
                'record': record,
                'error': str(e),
                'timestamp': datetime.now().isoformat()
            })
            logging.warning(f"记录清洗失败 id={record.get('id')}: {e}")

    return success, failed
第四步:批量写入(事务控制)

代码语言:javascript
复制
def batch_insert(pool, table_name, records, columns):
    """
    使用executemany批量插入,配合事务保证原子性
    """
    if not records:
        return 0

    conn = pool.get_connection()
    cursor = conn.cursor()

    # 构造INSERT语句
    placeholders = ', '.join(['%s'] * len(columns))
    sql = f"INSERT INTO {table_name} ({', '.join(columns)}) VALUES ({placeholders})"

    try:
        # 提取值列表
        values = [[r[col] for col in columns] for r in records]

        cursor.executemany(sql, values)
        conn.commit()

        return cursor.rowcount

    except Exception as e:
        conn.rollback()
        logging.error(f"批量插入失败: {e}")
        # 降级为单条插入,进一步隔离异常
        success_count = 0
        for record in records:
            try:
                single_values = [record[col] for col in columns]
                cursor.execute(sql, single_values)
                conn.commit()
                success_count += 1
            except Exception as e2:
                logging.error(f"单条插入失败 id={record.get('id')}: {e2}")
        return success_count

    finally:
        cursor.close()
        conn.close()
第五步:主控流程(断点续跑)

代码语言:javascript
复制
import json
import os
from datetime import datetime

CHECKPOINT_FILE = 'migrate_checkpoint.json'

def load_checkpoint():
    """加载断点,支持任务中断后从上次位置继续"""
    if os.path.exists(CHECKPOINT_FILE):
        with open(CHECKPOINT_FILE, 'r', encoding='utf-8') as f:
            return json.load(f)
    return {'last_processed_id': 0, 'total_success': 0, 'total_failed': 0}

def save_checkpoint(last_id, success_count, failed_count):
    """保存断点"""
    with open(CHECKPOINT_FILE, 'w', encoding='utf-8') as f:
        json.dump({
            'last_processed_id': last_id,
            'total_success': success_count,
            'total_failed': failed_count,
            'updated_at': datetime.now().isoformat()
        }, f, ensure_ascii=False, indent=2)

def main():
    pool = create_connection_pool()
    checkpoint = load_checkpoint()

    total_success = checkpoint['total_success']
    total_failed = checkpoint['total_failed']
    last_id = checkpoint['last_processed_id']

    logging.info(f"从断点恢复: last_id={last_id}, 已处理成功={total_success}")

    batch_count = 0
    for batch in fetch_in_batches(pool, batch_size=5000):
        batch_count += 1

        # 跳过已处理的批次(简单实现:基于ID过滤)
        if batch and batch[0]['id'] <= last_id:
            continue

        # 清洗
        success_records, failed_records = process_batch(batch)

        # 写入
        columns = ['id', 'order_no', 'phone', 'province', 'amount', 'create_time']
        inserted = batch_insert(pool, 'orders_summary', success_records, columns)

        total_success += inserted
        total_failed += len(failed_records)
        last_id = batch[-1]['id']

        # 每10个批次保存一次断点
        if batch_count % 10 == 0:
            save_checkpoint(last_id, total_success, total_failed)
            logging.info(f"进度: 批次={batch_count}, 成功={total_success}, 失败={total_failed}")

    # 最终保存
    save_checkpoint(last_id, total_success, total_failed)
    logging.info(f"任务完成: 总计成功={total_success}, 失败={total_failed}")

    # 失败记录单独导出,供人工复核
    if failed_records:
        with open('failed_records.json', 'w', encoding='utf-8') as f:
            json.dump(failed_records, f, ensure_ascii=False, indent=2)

if __name__ == '__main__':
    main()

五、RPA工具对比:为什么流程自动化比脚本更适合批量处理

5.1 可视化编排 + 自定义界面,降低维护成本

上面那段代码,如果交给一个不懂Python的运营同事维护,基本等于"天书"。但在流程自动化软件里,整个逻辑可以拖拽成一张流程图,还能自定义界面,拖拽组件设计出属于自己的软件界面,让非技术人员也能看懂流程在做什么:

代码语言:javascript
复制
[定时触发/API触发] → [连接数据库] → [分批查询] → [数据清洗] → [批量写入] → [发送邮件]
         ↓                  ↓              ↓              ↓              ↓
    [定时任务调度]     [异常重试]    [脏数据隔离]    [事务提交]    [生成报表]

每个节点独立配置,出了问题一眼就能定位到哪个环节。而且支持API触发,可以无缝对接现有的业务系统;也支持定时执行,设定好时间规则就能自动跑,不用人工守着。

5.2 内置容错机制,RPA比脚本稳定得多

脚本需要自己写try-except、写重试逻辑、写日志、写断点。流程自动化工具把这些能力内置了:

  • 异常捕获:某个节点失败,自动重试3次,仍失败则走异常分支
  • 断点续跑:任务中断后,从上次成功位置继续,不用从头再来
  • 脏数据隔离:异常记录自动进入"待复核队列",不影响主流程

这就是为什么RPA比脚本更稳——不是脚本做不到,而是脚本需要你一行行写,漏一个就崩。而自动化工具把这些工程能力封装好了,开箱即用。

5.3 跨系统协同,天生擅长

百万级数据处理很少是"单机游戏"。流程自动化工具天然支持:

  • 从MySQL读取 → 清洗 → 写入PostgreSQL
  • 处理完成后自动调用钉钉/飞书/企微机器人发通知
  • 生成Excel/PDF报表并邮件发送
  • 对接指纹浏览器做网页端数据校验(比如紫鸟浏览器、比特浏览器、hubstudio浏览器、adspower浏览器等)

这些如果用纯脚本实现,光是各个SDK的对接就要写几百行代码。而且RPA工具还能支持脚本打包导出EXE,把整套流程封装成一个可执行文件,发给同事双击就能跑,多设备使用无需多开会员

5.4 离线运行,数据不出本地

很多企业(尤其是金融、政务、医疗行业)对数据安全要求极高,数据不能离开本地服务器。一些内网自动化工具支持完全离线部署,内网离线使用数据不出本地,所有流程应用数据全部保存在用户本地设备上,不同步到服务端。这对于处理敏感数据来说,是脚本方案难以比拟的优势。

5.5 打包分发 + 在线更新 + 授权控制

脚本方案最大的痛点之一是环境依赖。Python版本、第三方库、数据库驱动,换台机器跑不起来是常事。而优秀的RPA软件支持将流程打包成独立的EXE应用,发给同事双击就能运行,不需要安装任何客户端或配置环境。还能给EXE设置授权,控制谁能用、用多久——这就是打包导出应用EXE支持授权的能力。

更实用的是,打包后的EXE支持在线推送更新,你修复了一个Bug或者优化了流程逻辑,不需要再手动一个个发文件,只需打开应用就能自动检测更新新版本。如果需要团队协作,还能通过加密分享的方式把应用安全地分发给指定人员,配合分享授权机制,做到"谁可以用、用到什么时候"完全可控。


六、性能对比:脚本 vs 流程自动化

维度

传统脚本

流程自动化方案

开发时间

2-3天(含调试)

2-4小时(拖拽配置)

异常处理

需手写,易遗漏

内置,开箱即用

断点续跑

需自行实现

原生支持

跨系统对接

需写胶水代码

可视化配置

维护成本

高(代码难读)

低(流程图直观)

部署方式

依赖环境

可打包EXE独立运行

数据安全

依赖脚本所在环境

可完全离线,数据不出本地

百万级耗时

约45分钟(含调试)

约30分钟(稳定运行)

注:以上数据基于MySQL 8.0、单表120万条记录、本地SSD的测试环境。实际耗时受硬件、网络、数据复杂度影响。

从这张表能看出,RPA比脚本快不只是执行速度,更重要的是开发速度和稳定性。RPA性能优化的核心不是让单次执行快多少,而是让整个流程从开发到部署到维护的全链路效率提升。


七、进阶:AI+RPA让数据处理更智能

7.1 AI辅助元素定位:不用学xpath

在处理Web端数据时(比如从后台管理系统获取数据再入库),经常遇到页面元素变动导致脚本失效的问题。现在一些先进的流程自动化软件已经接入了大模型能力,可以用自然语言描述元素,AI自动生成稳定的XPath路径——你不需要再去啃那些晦涩难懂的xpath语法,直接说"页面右上角的提交按钮",AI就能帮你定位。

这背后是AI智能优化元素路径的能力:元素获取支持本地智能生成,可根据生成结果选择合适稳定的元素路径。AI智能优化元素优化元素路径,无需学习晦涩难懂的xpath语法,通过自然语言描述即生成对应的xpath路径。甚至当Web元素失效时,AI能自动修复元素定位,实现"元素自愈",保障流程不中断。

7.2 Agent智能调度:在微信里指挥流程

更进一步,可以通过Agent功能,在钉钉、飞书、企微、个人微信里直接发送指令控制流程执行。比如@机器人说"跑一下昨天的订单统计",流程自动触发,完成后把结果回调到群里。这种对话式调度让非技术人员也能轻松使用自动化能力。

这里用的是最新的DeepseekV4模型做智能指令解析,回调通知响应执行结果等操作都能自动完成。

7.3 大模型赋能数据清洗:告别正则表达式

接入文心一言、豆包、DeepSeek、Kimi等大模型后,AI数据清洗不再依赖死板的正则表达式。比如让AI判断一条地址记录是否规范、自动补全省市信息、识别并修正错别字。甚至支持图片识图与OCR功能,直接把截图里的表格数据提取出来入库。

费用方面,AI功能采用用户自行对接各平台API的方式,用多少付多少,费用更可控。不像某些工具把AI功能打包成增值服务按月收费,这种"自带API Key"的模式对中小企业个人工作室更友好。


八、工具选型的一点建议

百万级数据处理,技术本身不难,难的是工程化——如何让流程稳定、可维护、可交接、可扩展。

脚本就像一把刀,灵活但容易伤到自己。流程自动化工具更像一台数控机床,前期配置略花时间,但一旦跑起来,稳定性、可观测性、可维护性都远超脚本。

对于个人开发者工作室或者中小企业来说,选择一款免费版使用无使用时长限制无运行时长无流程数量限制支持打包EXE分发可离线运行能对接主流大模型RPA软件,可能是比"从头造轮子"更务实的选择。

毕竟,我们的目标不是写出最漂亮的代码,而是让数据准时、准确、安全地到达目的地。

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

评论
登录后参与评论
0 条评论
热度
最新
推荐阅读
目录
  • 一、当百万级数据成为日常
  • 二、百万级数据处理的三大痛点
    • 2.1 连接池撑不住:一次性加载百万条的代价
    • 2.2 异常中断无恢复:一条脏数据毁掉整批任务
    • 2.3 跨系统协同困难:脚本不是万能的
  • 三、RPA+数据库:一种更工程化的自动化数据处理方案
    • 3.1 什么是RPA+数据库方案
    • 3.2 架构设计:分层解耦
  • 四、实战:百万订单数据迁移与清洗
    • 4.1 场景描述
    • 4.2 传统脚本方案(及其隐患)
    • 4.3 流程自动化方案:分批+容错+自愈
      • 第一步:数据库连接配置
      • 第二步:分批读取(游标方式)
      • 第三步:数据清洗(带异常隔离)
      • 第四步:批量写入(事务控制)
      • 第五步:主控流程(断点续跑)
  • 五、RPA工具对比:为什么流程自动化比脚本更适合批量处理
    • 5.1 可视化编排 + 自定义界面,降低维护成本
    • 5.2 内置容错机制,RPA比脚本稳定得多
    • 5.3 跨系统协同,天生擅长
    • 5.4 离线运行,数据不出本地
    • 5.5 打包分发 + 在线更新 + 授权控制
  • 六、性能对比:脚本 vs 流程自动化
  • 七、进阶:AI+RPA让数据处理更智能
    • 7.1 AI辅助元素定位:不用学xpath
    • 7.2 Agent智能调度:在微信里指挥流程
    • 7.3 大模型赋能数据清洗:告别正则表达式
  • 八、工具选型的一点建议
领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档