使用boto3获取S3桶多子文件夹最新文件并写入SQLite3数据库
解决S3子文件夹最新JSON文件获取及SQLite写入问题
一、修复S3子文件夹最新文件获取问题
你的代码只拿到单个文件夹文件、日期错误,核心原因是没按S3的前缀(模拟文件夹)分组筛选,且可能错误处理了文件修改时间。以下是修正后的实现:
核心思路
S3没有真实文件夹,所有对象通过key的前缀模拟层级。我们需要:
- 遍历桶内所有JSON文件
- 按
key的前缀(子文件夹路径)分组 - 每组内筛选
last_modified最新的文件
代码实现
import boto3 from collections import defaultdict import json import sqlite3 from datetime import datetime # 初始化S3资源(确保本地已配置AWS凭证,或通过参数传入) s3 = boto3.resource('s3') bucket_name = 'your-target-bucket' # 替换为你的桶名 bucket = s3.Bucket(bucket_name) # 存储每个子文件夹的最新文件信息 folder_latest_files = defaultdict(lambda: {'last_modified': datetime.min, 'key': None}) # 遍历桶内所有对象,筛选JSON文件并分组 for obj in bucket.objects.all(): # 跳过S3自动生成的"文件夹标记"对象(key以/结尾) if obj.key.endswith('/'): continue # 只处理JSON文件 if not obj.key.endswith('.json'): continue # 提取子文件夹路径:例如key为"sub-folder1/sub-folder2/file.json",则文件夹路径为"sub-folder1/sub-folder2/" folder_parts = obj.key.split('/')[:-1] folder_path = '/'.join(folder_parts) + '/' if folder_parts else '' # 比较并更新当前文件夹的最新文件 if obj.last_modified > folder_latest_files[folder_path]['last_modified']: folder_latest_files[folder_path] = { 'last_modified': obj.last_modified, 'key': obj.key } # 提取所有子文件夹的最新文件key列表 latest_json_keys = [info['key'] for info in folder_latest_files.values() if info['key']]
关键说明
obj.last_modified是原生datetime对象,直接比较即可,避免手动转换字符串导致的错误- 用
defaultdict自动初始化每个文件夹的默认值,简化分组逻辑 - 跳过以
/结尾的对象:这是S3用来标记文件夹存在的虚拟对象,不是实际文件
二、JSON解析并写入SQLite3
完成最新文件获取后,只需读取S3文件内容、解析JSON,再写入SQLite即可。以下是完整流程:
1. 初始化SQLite并创建表
根据你的JSON结构调整表字段,示例中假设JSON包含user_id、content、create_time三个字段:
# 初始化SQLite连接(文件不存在则自动创建) conn = sqlite3.connect('s3_json_db.sqlite') cursor = conn.cursor() # 创建数据表(按需调整字段) cursor.execute(''' CREATE TABLE IF NOT EXISTS s3_json_records ( id INTEGER PRIMARY KEY AUTOINCREMENT, folder_path TEXT, s3_key TEXT UNIQUE, file_last_modified TEXT, user_id INTEGER, content TEXT, create_time TEXT ) ''') conn.commit()
2. 读取S3文件并写入数据库
# 遍历每个最新文件,完成读取-解析-写入 for folder_path, file_info in folder_latest_files.items(): s3_key = file_info['key'] if not s3_key: continue # 读取S3中的JSON文件 s3_obj = s3.Object(bucket_name, s3_key) json_str = s3_obj.get()['Body'].read().decode('utf-8') try: json_data = json.loads(json_str) except json.JSONDecodeError as e: print(f"解析文件 {s3_key} 失败: {str(e)}") continue # 格式化修改时间为字符串 modified_time_str = file_info['last_modified'].strftime('%Y-%m-%d %H:%M:%S') # 插入数据(用INSERT OR REPLACE避免重复写入同一文件) try: cursor.execute(''' INSERT OR REPLACE INTO s3_json_records (folder_path, s3_key, file_last_modified, user_id, content, create_time) VALUES (?, ?, ?, ?, ?, ?) ''', ( folder_path, s3_key, modified_time_str, json_data.get('user_id'), json_data.get('content'), json_data.get('create_time') )) conn.commit() except Exception as e: print(f"写入文件 {s3_key} 失败: {str(e)}") conn.rollback() # 关闭数据库连接 conn.close()
关键说明
- 使用
INSERT OR REPLACE:因为s3_key设为UNIQUE,如果同一文件更新,会自动替换旧数据 - 解析JSON时加入异常捕获:避免单个文件格式错误导致整个流程中断
- 按需调整表字段:如果你的JSON结构复杂,也可以直接存储整个JSON字符串(用
json.dumps(json_data)),后续查询时再解析
完整整合代码
将上述两部分代码合并,替换bucket_name和表字段即可直接运行。
内容的提问来源于stack exchange,提问作者KLC2021
相关产品推荐
相关产品推荐

