Python脚本如何触发增量数据同步任务

wen python案例 29

本文目录导读:

Python脚本如何触发增量数据同步任务

  1. 时间戳机制
  2. 版本号/序列号机制
  3. 状态标记触发
  4. 基于事件/消息队列触发
  5. 定时任务触发
  6. 完整示例:API触发增量同步
  7. 选择建议

我来介绍几种在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)

选择建议

  1. 数据量小、实时性要求低:使用时间戳机制
  2. 需要精确控制同步顺序:使用版本号机制
  3. 需要重试机制:使用状态标记
  4. 高并发、实时性强:使用消息队列
  5. 定期执行:使用定时任务
  6. 需要外部触发:使用API接口

选择合适的触发方式取决于你的具体业务场景和技术栈。

抱歉,评论功能暂时关闭!