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

使用boto3获取S3桶多子文件夹最新文件并写入SQLite3数据库

解决S3子文件夹最新JSON文件获取及SQLite写入问题

一、修复S3子文件夹最新文件获取问题

你的代码只拿到单个文件夹文件、日期错误,核心原因是没按S3的前缀(模拟文件夹)分组筛选,且可能错误处理了文件修改时间。以下是修正后的实现:

核心思路

S3没有真实文件夹,所有对象通过key的前缀模拟层级。我们需要:

  1. 遍历桶内所有JSON文件
  2. 按key的前缀(子文件夹路径)分组
  3. 每组内筛选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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 20:54:32