本文目录导读:

我来介绍几种在Python脚本中触发增量数据同步任务的方法:
时间戳机制
import datetime
import pandas as pd
from sqlalchemy import create_engine, text
class IncrementalSync:
def __init__(self, connection_string):
self.engine = create_engine(connection_string)
self.last_sync_time = None
self.sync_config_file = 'sync_config.json'
def load_last_sync_time(self):
"""加载上次同步时间"""
try:
import json
with open(self.sync_config_file, 'r') as f:
config = json.load(f)
self.last_sync_time = datetime.datetime.fromisoformat(config['last_sync_time'])
except:
self.last_sync_time = datetime.datetime(2024, 1, 1) # 默认起始时间
def save_sync_time(self):
"""保存本次同步时间"""
import json
config = {'last_sync_time': self.current_sync_time.isoformat()}
with open(self.sync_config_file, 'w') as f:
json.dump(config, f)
def get_incremental_data(self, table_name, timestamp_column='updated_at'):
"""获取增量数据"""
self.load_last_sync_time()
self.current_sync_time = datetime.datetime.now()
query = f"""
SELECT * FROM {table_name}
WHERE {timestamp_column} > :last_sync
AND {timestamp_column} <= :current_sync
"""
df = pd.read_sql(query, self.engine, params={
'last_sync': self.last_sync_time,
'current_sync': self.current_sync_time
})
self.save_sync_time()
return df
# 使用示例
sync = IncrementalSync('postgresql://user:pass@localhost:5432/db')
df = sync.get_incremental_data('orders')
print(f"获取到 {len(df)} 条增量数据")
版本号/序列号机制
class VersionSync:
def __init__(self, db_connection):
self.connection = db_connection
self.last_version = 0
def get_max_version(self, table_name, version_column='version'):
"""获取当前最大版本号"""
query = f"SELECT MAX({version_column}) FROM {table_name}"
cursor = self.connection.cursor()
cursor.execute(query)
result = cursor.fetchone()[0]
return result if result else 0
def sync_data(self, table_name, version_column='version', batch_size=1000):
"""按版本号同步数据"""
max_version = self.get_max_version(table_name, version_column)
current_version = self.last_version
while current_version < max_version:
next_version = min(current_version + batch_size, max_version)
query = f"""
SELECT * FROM {table_name}
WHERE {version_column} > {current_version}
AND {version_column} <= {next_version}
ORDER BY {version_column}
"""
# 执行同步逻辑
self.process_batch(query)
self.last_version = next_version
current_version = next_version
def process_batch(self, query):
"""处理批次数据"""
# 实现具体的数据处理逻辑
pass
状态标记触发
class StatusTriggerSync:
def __init__(self, source_db, target_db):
self.source_db = source_db
self.target_db = target_db
self.sync_status_table = 'sync_status'
def mark_records_for_sync(self, table_name, ids_to_sync):
"""标记需要同步的记录"""
import mysql.connector
cursor = self.source_db.cursor()
update_query = f"""
UPDATE {table_name}
SET sync_status = 'pending',
sync_attempts = IFNULL(sync_attempts, 0) + 1
WHERE id IN ({','.join(map(str, ids_to_sync))})
"""
cursor.execute(update_query)
self.source_db.commit()
def trigger_sync_job(self):
"""触发同步作业"""
cursor = self.source_db.cursor()
# 获取待同步记录
query = f"""
SELECT * FROM {self.source_table}
WHERE sync_status = 'pending'
AND sync_attempts < 3
"""
cursor.execute(query)
records = cursor.fetchall()
# 执行同步
for record in records:
try:
# 同步到目标数据库
self.sync_to_target(record)
# 更新状态为已同步
update_query = f"""
UPDATE {self.source_table}
SET sync_status = 'completed',
sync_time = NOW()
WHERE id = {record['id']}
"""
cursor.execute(update_query)
except Exception as e:
# 同步失败,更新错误信息
update_query = f"""
UPDATE {self.source_table}
SET sync_status = 'failed',
error_message = '{e}',
last_error_time = NOW()
WHERE id = {record['id']}
"""
cursor.execute(update_query)
self.source_db.commit()
# 使用示例
sync = StatusTriggerSync(source_db, target_db)
sync.trigger_sync_job() # 触发增量同步
基于事件/消息队列触发
import redis
import json
from kafka import KafkaConsumer, KafkaProducer
class EventDrivenSync:
def __init__(self):
self.redis_client = redis.Redis(host='localhost', port=6379, db=0)
self.kafka_producer = KafkaProducer(
bootstrap_servers=['localhost:9092'],
value_serializer=lambda v: json.dumps(v).encode('utf-8')
)
def publish_change_event(self, table_name, operation, data):
"""发布数据变更事件"""
event = {
'table': table_name,
'operation': operation, # 'INSERT', 'UPDATE', 'DELETE'
'data': data,
'timestamp': '2024-01-15 10:30:00'
}
# 发布到Kafka
self.kafka_producer.send('data_changes', value=event)
# 同时也发布到Redis
self.redis_client.publish('data_changes', json.dumps(event))
def consume_events(self):
"""消费事件并执行同步"""
consumer = KafkaConsumer(
'data_changes',
bootstrap_servers=['localhost:9092'],
group_id='sync_group',
value_deserializer=lambda m: json.loads(m.decode('utf-8'))
)
for message in consumer:
event = message.value
self.process_event(event)
def process_event(self, event):
"""处理单个事件"""
table = event['table']
operation = event['operation']
data = event['data']
if operation == 'INSERT':
# 执行插入同步
self.sync_insert(table, data)
elif operation == 'UPDATE':
# 执行更新同步
self.sync_update(table, data)
elif operation == 'DELETE':
# 执行删除同步
self.sync_delete(table, data)
# Webhook触发示例
@app.route('/webhook/data-change', methods=['POST'])
def handle_data_change():
"""Webhook接收数据变更通知"""
data = request.json
sync = EventDrivenSync()
sync.publish_change_event(
table_name=data['table'],
operation=data['operation'],
data=data['record']
)
return jsonify({'status': 'success'})
定时任务触发
import schedule
import time
from datetime import datetime
class ScheduledSync:
def __init__(self):
self.sync_tasks = {}
def add_sync_task(self, name, sync_func, interval_minutes=5):
"""添加定时同步任务"""
self.sync_tasks[name] = sync_func
schedule.every(interval_minutes).minutes.do(self.execute_task, name)
def execute_task(self, task_name):
"""执行同步任务"""
print(f"[{datetime.now()}] 执行同步任务: {task_name}")
try:
if task_name in self.sync_tasks:
self.sync_tasks[task_name]()
print(f"[{datetime.now()}] 同步任务完成: {task_name}")
except Exception as e:
print(f"[{datetime.now()}] 同步任务失败: {task_name}, 错误: {e}")
def run(self):
"""运行调度器"""
print("开始调度同步任务...")
# 立即执行一次所有任务
for task_name in self.sync_tasks:
self.execute_task(task_name)
# 持续调度
while True:
schedule.run_pending()
time.sleep(1)
# 使用示例
scheduler = ScheduledSync()
def sync_orders():
"""同步订单数据"""
# 实现订单增量同步逻辑
pass
def sync_users():
"""同步用户数据"""
# 实现用户增量同步逻辑
pass
# 每5分钟同步订单
scheduler.add_sync_task('orders', sync_orders, interval_minutes=5)
# 每1小时同步用户
scheduler.add_sync_task('users', sync_users, interval_minutes=60)
scheduler.run()
完整示例:API触发增量同步
from flask import Flask, request, jsonify
import threading
from datetime import datetime
app = Flask(__name__)
class IncrementalSyncManager:
def __init__(self):
self.sync_lock = threading.Lock()
self.sync_status = {
'last_sync_time': None,
'is_syncing': False,
'total_records_synced': 0,
'errors': []
}
def trigger_sync(self, sync_type='incremental'):
"""
触发增量同步
sync_type: 'incremental', 'full', 'specific_table'
"""
if self.sync_status['is_syncing']:
return {'status': 'error', 'message': '同步正在进行中'}
with self.sync_lock:
self.sync_status['is_syncing'] = True
self.sync_status['start_time'] = datetime.now()
# 在后台线程执行同步
thread = threading.Thread(target=self._execute_sync, args=(sync_type,))
thread.start()
return {'status': 'success', 'message': '同步任务已启动'}
def _execute_sync(self, sync_type):
"""执行同步逻辑"""
try:
# 获取增量数据
if sync_type == 'incremental':
data = self.get_incremental_data()
elif sync_type == 'full':
data = self.get_full_data()
else:
data = self.get_table_data(sync_type)
# 处理同步
synced_count = self.process_sync(data)
# 更新状态
self.sync_status['total_records_synced'] += synced_count
self.sync_status['last_sync_time'] = datetime.now()
except Exception as e:
self.sync_status['errors'].append({
'time': datetime.now(),
'error': str(e)
})
finally:
self.sync_status['is_syncing'] = False
self.sync_status['end_time'] = datetime.now()
sync_manager = IncrementalSyncManager()
@app.route('/api/sync/trigger', methods=['POST'])
def trigger_sync():
"""API触发同步"""
data = request.json
sync_type = data.get('sync_type', 'incremental')
result = sync_manager.trigger_sync(sync_type)
return jsonify(result)
@app.route('/api/sync/status', methods=['GET'])
def get_sync_status():
"""获取同步状态"""
return jsonify(sync_manager.sync_status)
@app.route('/api/sync/schedule', methods=['POST'])
def schedule_sync():
"""设置定时同步"""
data = request.json
interval_minutes = data.get('interval_minutes', 30)
sync_type = data.get('sync_type', 'incremental')
# 设置定时任务
schedule.every(interval_minutes).minutes.do(
sync_manager.trigger_sync, sync_type
)
return jsonify({
'status': 'success',
'message': f'已设置每{interval_minutes}分钟执行{sync_type}同步'
})
if __name__ == '__main__':
# 启动调度器线程
import schedule
def run_scheduler():
while True:
schedule.run_pending()
time.sleep(1)
threading.Thread(target=run_scheduler, daemon=True).start()
app.run(host='0.0.0.0', port=5000)
选择建议
- 数据量小、实时性要求低:使用时间戳机制
- 需要精确控制同步顺序:使用版本号机制
- 需要重试机制:使用状态标记
- 高并发、实时性强:使用消息队列
- 定期执行:使用定时任务
- 需要外部触发:使用API接口
选择合适的触发方式取决于你的具体业务场景和技术栈。