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

如何用Python实现每日将Azure Blob Storage文件无丢失同步至Azure SQL数据库

如何用Python实现每日将Azure Blob Storage文件无丢失同步至Azure SQL数据库

我懂你现在的难处——每天要把Azure Blob里的文件同步到SQL数据库,还得防着两边服务偶尔掉链子导致数据丢失,关键是Azure Data Factory对parquet文件的支持又不给力,只能靠Python自己撸方案对吧?结合你已经列出的那些库,我给你整理一套靠谱的无丢失同步方案,核心就是状态追踪+重试机制+幂等性保障,咱们一步步来:

一、核心设计思路(无丢失的关键)

  • 状态追踪:记录每个Blob文件的同步状态(未同步、同步中、已完成),避免重复同步或者漏同步
  • 幂等性保障:即使重复运行脚本,也不会把同一份数据多次写入SQL
  • 失败重试:遇到Blob或SQL的临时故障时自动重试,减少人工干预
  • 异常日志:详细记录每一步的错误,方便排查问题

二、完整实现代码

下面的代码完全基于你已经导入的库来写,直接就能复用:

1. 初始化配置与日志

先加载环境变量、配置日志,出问题时能快速定位:

load_dotenv()  # 从.env文件加载Azure认证信息、SQL连接串等

# 配置日志,同时输出到文件和控制台
logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s - %(levelname)s - %(message)s',
    handlers=[logging.FileHandler('sync_log.log'), logging.StreamHandler(sys.stdout)]
)
logger = logging.getLogger(__name__)

2. 认证与客户端初始化

用你已经用到的ClientSecretCredential完成Blob和SQL的身份认证:

# 初始化Azure Blob客户端
credential = ClientSecretCredential(
    tenant_id=os.getenv("AZURE_TENANT_ID"),
    client_id=os.getenv("AZURE_CLIENT_ID"),
    client_secret=os.getenv("AZURE_CLIENT_SECRET")
)

blob_service_client = BlobServiceClient(
    account_url=f"https://{os.getenv('AZURE_STORAGE_ACCOUNT')}.blob.core.windows.net",
    credential=credential
)
container_client = blob_service_client.get_container_client(os.getenv("BLOB_CONTAINER_NAME"))

# 初始化Azure SQL连接(用SQLAlchemy简化操作)
params = urllib.parse.quote_plus(os.getenv("AZURE_SQL_CONNECTION_STRING"))
engine = create_engine(f"mssql+pyodbc:///?odbc_connect={params}")

3. 状态记录机制(核心!避免丢失)

咱们用本地pickle文件记录同步状态(怕本地文件丢的话,也可以改成存在SQL的一张小表):

STATE_FILE = "sync_state.pkl"

def load_sync_state():
    """加载已同步的Blob文件状态"""
    if os.path.exists(STATE_FILE):
        with open(STATE_FILE, 'rb') as f:
            return pickle.load(f)
    # 初始状态:空的已完成集合和同步中集合
    return {"completed_blobs": set(), "in_progress": set()}

def save_sync_state(state):
    """保存同步状态,防止脚本中断后丢失进度"""
    with open(STATE_FILE, 'wb') as f:
        pickle.dump(state, f)

4. 同步逻辑(带重试与幂等)

针对parquet文件读取和SQL写入做了重试处理,同时保证不会重复写入:

def sync_blob_to_sql(blob_name, max_retries=3):
    """单个Blob文件同步到SQL,带指数退避重试"""
    retry_count = 0
    while retry_count < max_retries:
        try:
            # 1. 下载Blob内容到内存
            blob_client = container_client.get_blob_client(blob_name)
            blob_data = blob_client.download_blob().readall()
            
            # 2. 读取Parquet文件(用polars比pandas更快,也可以替换成pandas)
            df = pl.read_parquet(io.BytesIO(blob_data))
            
            # 3. 写入SQL(保证幂等性:先检查是否已同步过)
            with engine.connect() as conn:
                # 假设你的SQL表有blob_name字段用来标识唯一文件
                exists = conn.execute(
                    sqlalchemy.text(f"SELECT 1 FROM your_target_table WHERE blob_name = :name"), 
                    {"name": blob_name}
                ).fetchone()
                if not exists:
                    # 给数据加上Blob文件名标识,方便后续校验
                    df_with_blob = df.with_columns(pl.lit(blob_name).alias("blob_name"))
                    df_with_blob.write_database(
                        table_name="your_target_table",
                        connection=engine,
                        if_exists="append"
                    )
                    logger.info(f"✅ 成功同步Blob: {blob_name}")
                else:
                    logger.info(f"⏭️ Blob {blob_name}已同步过,跳过")
            return True
        except (ResourceNotFoundError, ClientAuthenticationError, sqlalchemy.exc.OperationalError) as e:
            retry_count += 1
            wait_time = 2 ** retry_count  # 指数退避:2s、4s、8s...
            logger.error(f"❌ 同步Blob {blob_name}失败,{wait_time}s后重试第{retry_count}次: {str(e)}")
            time.sleep(wait_time)
        except Exception as e:
            logger.error(f"💥 同步Blob {blob_name}遇到未知错误,不再重试: {str(e)}")
            return False
    logger.error(f"❌ 同步Blob {blob_name}超过最大重试次数,失败")
    return False

def daily_sync():
    """每日同步主逻辑"""
    state = load_sync_state()
    completed_blobs = state["completed_blobs"]
    in_progress = state["in_progress"]
    
    # 获取容器内今天的Blob文件(假设Blob名包含日期,比如"2024-05-20_data.parquet")
    today = datetime.now().strftime("%Y-%m-%d")
    blobs = container_client.list_blobs()
    target_blobs = [
        b.name for b in blobs 
        if today in b.name 
        and b.name not in completed_blobs 
        and b.name not in in_progress
    ]
    
    logger.info(f"🔍 找到待同步Blob数量: {len(target_blobs)}")
    
    for blob_name in tqdm(target_blobs, desc="同步进度"):
        # 标记为同步中,防止脚本中断后重复处理
        in_progress.add(blob_name)
        save_sync_state(state)
        
        # 执行同步
        success = sync_blob_to_sql(blob_name)
        
        # 更新状态
        if success:
            completed_blobs.add(blob_name)
        in_progress.remove(blob_name)
        save_sync_state(state)
    
    logger.info("🎉 每日同步任务完成")

5. 定时自动运行

要实现每日自动执行,可以用schedule库,或者系统的定时任务(Linux用crontab,Windows用任务计划程序):

import schedule

# 设置每天凌晨2点运行同步任务
schedule.every().day.at("02:00").do(daily_sync)

logger.info("⏰ 定时同步服务启动,等待任务执行...")
while True:
    schedule.run_pending()
    time.sleep(60)  # 每分钟检查一次任务

三、额外的无丢失保障建议

  • 状态持久化升级:如果担心本地pickle文件损坏,可以把同步状态存在Azure SQL的一张专门的小表(比如sync_status),包含blob_name、sync_status、sync_time字段
  • 数据校验:同步完成后,可以对比Blob的行数和SQL里的行数,确保数据完全一致
  • 异常告警:可以加邮件/Teams机器人告警,当同步失败次数过多时及时通知你
  • 表结构兼容:如果parquet文件结构有变化,建议在写入前先检查表结构,避免写入失败

备注:内容来源于stack exchange,提问作者ondoondo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 18:28:03