基于Lambda脚本的MongoDB唯一数据插入方案及问题咨询
问题:基于复合主键的MongoDB增量插入实现疑问
需求
通过Lambda脚本向MongoDB插入数据,仅当数据基于复合主键(item_id、process_name、actor)不存在时执行插入。
当前实现步骤
- 从CSV加载数据到DataFrame并保留指定表头:
df = wr.s3.read_csv(path=part_file_url, low_memory=False) df = df[header]
- 查询MongoDB集合,仅获取复合主键字段数据到
curr_data_df:
project_keys ={'_id': 0, 'item_id': 1, 'process_name': 1, 'actor': 1} queryString = [ "item_id", "process_name","actor"] curr_data_df = pd.DataFrame(list(audit_collection.find({"process_name":process_name, "processId":process_id}, project_keys)), columns=queryString)
- 通过左连接合并筛选唯一待插入数据:
result = df.merge(curr_data_df,how='left', indicator=True) #step 3.a result = result[result['_merge']=='left_only'] # step 3.b df = result.drop(columns='_merge') # step 3.c
获取唯一数据后执行批量插入。
咨询问题
- 上述基于复合主键左连接筛选唯一数据的方案是否正确?
- 合并时未指定连接列,是否会默认使用两个DataFrame共有的
item_id、process_name、actor作为连接列? - 重复上传同一文件时,原df有148条记录,curr_data_df有684条,但合并后结果有4166条记录,最终筛选出0条新数据,为何合并记录数异常?
解答
- 方案核心思路是可行的:通过左连接定位CSV中MongoDB不存在的记录,仅插入新数据。但存在潜在风险——如果MongoDB数据量较大,将所有复合主键拉取到DataFrame会占用大量内存,在Lambda这类资源受限环境中可能触发内存溢出。后续可考虑改用MongoDB的
bulkWrite结合updateOne(upsert: false),或通过聚合查询前置对比,减少本地数据处理量。 - 是的,
pd.merge默认会以两个DataFrame中名称完全一致的列作为连接键。但建议显式指定on=['item_id', 'process_name', 'actor'],避免后续新增同名列导致连接逻辑意外变更,同时提升代码可读性。 - 合并记录数异常主要有两类原因:
- 字段类型不匹配:比如
df中item_id是字符串类型,而curr_data_df中item_id是整数类型,看似相同的值因类型差异无法匹配,触发局部笛卡尔积,导致合并记录数暴增。 - 存在重复主键记录:若
curr_data_df中存在重复的(item_id, process_name, actor)组合,df中的单条记录会与多条重复记录匹配,生成多条合并结果;另外要检查MongoDB查询条件{"process_name":process_name, "processId":process_id}是否准确,是否拉取了无关的复合主键数据。
- 字段类型不匹配:比如
至于最终筛选出0条新数据,是重复上传的预期结果,但合并记录数异常是上述类型不匹配或重复数据问题导致的额外现象。
内容的提问来源于stack exchange,提问作者Abcd Efgh
相关产品推荐
相关产品推荐

