如何用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
相关产品推荐
相关产品推荐

