如何高效查询并存储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
相关产品推荐
相关产品推荐

