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

如何高效查询并存储5000万条Salesforce REST API记录构建数据湖

问题

问题场景

需要通过Salesforce REST API查询5000万+条记录,作为多源数据融合流程的一部分,后续用于机器学习模型训练。目标是将数据存储到数据库或Parquet文件中,支持后续操作。

Salesforce API返回机制

通过REST API查询大量对象数据时,Salesforce返回格式如下:

{"totalsize":50000000, "done":false, 
"nextrecordUrl":"/services/data/v49.0/a1000xxxxxx-2000",
 "records":[返回的2000条批量数据]}

其中nextRecordUrl用于获取下一批数据(Salesforce单批最大返回2000条,但数量不固定),当前通过循环遍历直至无nextRecordUrl获取全量数据。

当前实现方式

初始思路伪代码:

read the data = query(api)
while "done" is not True:
      read the "records" part of the json 
      load it into Database ( using psycopg2 )
      again read data = query( nextRecordUrl )
      read value of "done"

实际基于simple_salesforce库的实现代码:

sf = salesforce () # simple_salesforce library connection
df = pandas dataframe
sqlalc_engine = sqlalchemy engine
# query the first batch here then this loop start
while True:
            if num_batches == 1:
                current_df = df
            else:
                current_df = pd.DataFrame(lstRecords)
                #print(current_df.columns)
                current_df = current_df.drop(['attributes'], axis=1)

            print(f"next record URL is :{nextRecordsUrl}, || batch no is : {num_batches}")
            print(f"Loading batch : {num_batches} | records count : {current_df.shape[0]}")

            #this line i am using to insert records into DB
            current_df.to_sql(table_name, con=sqlalc_engine, schema='', if_exists='append', index=False)

            print(f"{num_batches} is loaded to DB")
            num_batches+=1

            if completed:
                #no more records
                break

            # get next batch of records
            records = sf.query_more(nextRecordsUrl, identifier_is_url=True)
            lstRecords = records.get('records')
            nextRecordsUrl = records.get('nextRecordsUrl')

            # set running check for next loop
            completed = records.get('done')

问题与诉求

当前流程即使单循环耗时2秒,全量数据入库也需要16-17小时,且需每两周执行一次同步最新数据。此外还面临数据库超时、内存溢出等问题,需要优化操作提升效率。


优化方案

1. 改用Salesforce Bulk API替代REST Query API

Salesforce Bulk API专为大规模数据导出设计,支持异步批量查询,单批次可处理远多于2000条的数据(最大支持10000条/批,甚至直接导出大文件),能大幅减少API请求次数。

示例实现(基于simple_salesforce):

from simple_salesforce import SalesforceBulk
import time

bulk = SalesforceBulk(username='xxx', password='xxx', security_token='xxx')
# 创建异步查询任务,指定导出格式为CSV
job = bulk.create_query_job('目标对象名', contentType='CSV')
# 提交查询语句
batch = bulk.query(job, "SELECT Id, Name, 字段1, 字段2 FROM 目标对象名")
# 关闭任务,触发执行
bulk.close_job(job)

# 轮询等待批次完成
while not bulk.is_batch_done(batch):
    time.sleep(10)

# 下载结果文件
result = bulk.get_batch_result(batch)

优势:减少API调用次数,降低网络开销;异步处理避免长时间同步阻塞;直接导出结构化文件,后续批量入库更高效。

2. 批量入库效率优化

当前单批to_sql插入效率极低,可通过以下方式优化:

  • 合并入库批次:累积5-10批Salesforce返回数据(1-2万条)后再执行一次插入,减少数据库连接和事务开销。
  • 使用数据库原生批量导入命令:如果目标是PostgreSQL,用psycopg2的copy_from替代to_sql,速度能提升10-100倍:
from io import StringIO
import pandas as pd

# 假设current_df是累积后的DataFrame
output = StringIO()
current_df.to_csv(output, sep='\t', header=False, index=False)
output.seek(0)

with sqlalc_engine.connect() as conn:
    cursor = conn.connection.cursor()
    cursor.copy_from(output, table_name, sep='\t', columns=current_df.columns)
    conn.commit()
  • 控制事务提交:在SQLAlchemy引擎中设置isolation_level='AUTOCOMMIT',或手动合并多个批次到一个事务中提交,减少事务次数。

3. 并行化处理(需注意API配额)

将数据获取和入库环节解耦并行,充分利用资源:

  • 并行查询:将查询按字段分段(比如按Id范围拆分:WHERE Id > 'xxx' AND Id < 'yyy'),用多线程同时发起Bulk API请求,最大化利用Salesforce的API速率配额。
  • 异步入库:用本地消息队列(如Redis)将获取到的数据放入队列,单独启动多个消费者进程执行入库操作,实现数据获取与入库的并行。

4. 优先选择Parquet文件存储(适配机器学习场景)

如果后续用于机器学习训练,直接存储为Parquet比入库数据库更高效:

  • Parquet是列存储格式,压缩率高、读写速度快,适合大数据量存储,且能直接被Spark、Dask等机器学习框架读取。
  • 示例代码:
current_df.to_parquet(f'output_batch_{num_batches}.parquet', engine='pyarrow')

后续可直接加载Parquet文件进行预处理,避免数据库导出的额外开销。

5. 增量同步替代全量同步

每两周全量同步5000万条数据效率极低,改为增量同步:

  • 记录每次同步的时间戳,下次同步时仅查询LastModifiedDate > '上次同步时间'的记录;优先使用SystemModstamp字段,比LastModifiedDate更准确。
  • 如果有自定义的增量标识字段(如同步标记),也可基于该字段筛选。

6. 内存与超时问题优化

  • 流式处理数据:边获取数据边写入文件/数据库,不要将全量数据加载到内存中。
  • 调整数据库连接配置:在SQLAlchemy引擎中增加超时参数:connect_args={'connect_timeout': 300},避免长时间等待超时。
  • 分块读取大文件:如果用Bulk API导出大文件,采用分块读取的方式处理,避免内存溢出。

内容的提问来源于stack exchange,提问作者Healer77Om

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 05:45:07