前言

Maxwell是一个开源的MySQL数据库binlog解析工具,用于将MySQL数据库的binlog转换成易于消费的JSON格式,并通过Kafka、RabbitMQ、Kinesis 等消息队列或直接写入文件等方式将其输出。本节内容主要介绍下如何在docker环境下以maxwell + redis + 自定义脚本的方式实现mysql 实时同步数据到pgsql。

一、mysql环境准备

1.1、登录mysql查看mysql有没有开启binlog日志:

SHOW VARIABLES LIKE 'log_bin';

1.2、如果没有开启则需要开启,辑 MySQL 配置文件 my.cnf,在 [mysqld] 部分添加或修改以下配置:

[mysqld]
# 启用 binlog
log_bin = /var/lib/mysql/mysql-bin

# 设置 binlog 格式,可选 STATEMENT、ROW、MIXED
binlog_format = ROW

# 设置 binlog 过期时间(可选,单位:天)
expire_logs_days = 7

# 设置每个 binlog 文件大小(可选,单位:字节)
max_binlog_size = 100M

# 确保每个事务都立即写入 binlog(可选)
sync_binlog = 1

1.3、重启 MySQL 服务,验证binlog是否开启

在这里插入图片描述

二、准备docker-compose文件,配置maxwell

2.1、准备docker-compose文件

version: '3.8'
services:
  pgsql:
    image: postgres:15-alpine
    container_name: pgsql
    restart: always
    environment:
      POSTGRES_DB: jkdb
      POSTGRES_USER: postgres
      POSTGRES_PASSWORD: 123456
    volumes:
      - ./pgsql-data:/var/lib/postgresql/data
 # 数据库初始化脚本
      - ./conf/init.sql:/docker-entrypoint-initdb.d/init.sql
    ports:
      - "5432:5432"
    networks:
      - maxwell-network

  redis:
    image: redis:5.0
    container_name: redis
    restart: always
    volumes:
      - ./conf/redis.conf:/usr/local/etc/redis/redis.conf
    command: redis-server /usr/local/etc/redis/redis.conf
    ports:
      - "6378:6379"
    networks:
      - maxwell-network

  maxwell:
    image: zendesk/maxwell:v1.40.0
    container_name: maxwell
    command: bin/maxwell --config /etc/maxwell/config.properties
    volumes:
      - ./conf:/etc/maxwell/
    depends_on:
      - pgsql
      - redis
    networks:
      - maxwell-network
  # 自定义消费者镜像,
  consumer:
    image: maxwell-python-consumer:latest
    container_name: consumer
    restart: always
    environment:
      PG_HOST: pgsql
      PG_PORT: 5432
      PG_USER: postgres
      PG_PASSWORD: 123456
      PG_DB: jkdb
      REDIS_HOST: redis
      REDIS_PORT: 6379
    depends_on:
      - redis
      - pgsql
    networks:
      - maxwell-network
networks:
  maxwell-network:
    driver: bridge

2.2、maxwell配置文件 config.properties

daemon=true
#info 级别记录一般信息,适合生产环境。还可设置为 debug、warn、error 等
log_level=info
# 用于区分多个 Maxwell 实例,尤其在向 Kafka 发送数据时,确保每个实例有独立标识
client_id=maxwell_1

# mysql连接相关配置
host=127.0.0.1
port=3306
# 连接mysql用户名,需要有监控其他数据库的权限
user=root
password=123456
#指定 Maxwell 存储元数据的数据库,程序启动自动创建,注意不是目标数据库
schema_database=maxwell
# 指定数据输出方式。
producer=redis
#redis连接相关配置
redis_host=redis
redis_port=6379
redis.db=0
# Redis数据写入类型,这里使用的是默认的pub/sub模式,可选值:xadd(Stream), list, channel, hash
redis_type=xadd

三、自定义消费者镜像

3.1 准备python脚本:这里只实现了基本的增删改,没有考虑语法兼容问题,如果要使用还需要完善语法兼容

import redis
import psycopg2
import json
import os

# 从环境变量中获取 Redis 配置
redis_host = os.getenv('REDIS_HOST', '10.29.4.196')
redis_port = int(os.getenv('REDIS_PORT', 6378))
redis_channel = os.getenv('REDIS_CHANNEL', 'maxwell')

# 从环境变量中获取 PostgreSQL 配置
pg_host = os.getenv('PG_HOST', '10.29.4.196')
pg_port = int(os.getenv('PG_PORT', 5432))
pg_user = os.getenv('PG_USER', 'postgres')
pg_password = os.getenv('PG_PASSWORD', '123456')
pg_database = os.getenv('PG_DATABASE', 'jkdb')

# 连接到 Redis
redis_client = redis.Redis(host=redis_host, port=redis_port)
# 订阅 Redis 频道
pubsub = redis_client.pubsub()
pubsub.subscribe(redis_channel)

# 连接到 PostgreSQL
try:
    conn = psycopg2.connect(
        host=pg_host,
        port=pg_port,
        user=pg_user,
        password=pg_password,
        database=pg_database
    )
    cursor = conn.cursor()
    print("成功连接到 PostgreSQL 数据库")
except psycopg2.Error as e:
    print(f"连接到 PostgreSQL 数据库时出错: {e}")
    exit(1)

# 开始接收和处理 Redis 消息
try:
    for message in pubsub.listen():
        if message['type'] == 'message':
            try:
                # 解析 JSON 消息
                data = json.loads(message['data'])
                print("接收到消息:")
                print(data)
                # 假设消息中有 'database'、'table' 和 'data' 字段
                database = data.get('database')
                table = data.get('table')
                row_data = data.get('data')

                if table and row_data:
                    # 构建插入语句
                    sql = ""

                    if data["type"] == "insert":
                        keys = ", ".join(data["data"].keys())
                        values = ", ".join([f"'{v}'" for v in data["data"].values()])
                        sql = f"INSERT INTO {database}.{table} ({keys}) VALUES ({values});"
                    elif data["type"] == "update":
                        updates = ", ".join([f"{k}='{v}'" for k, v in data["data"].items()])
                        sql = f"UPDATE {database}.{table} SET {updates} WHERE id={data['data']['id']};"
                    elif data["type"] == "delete":
                        sql = f"DELETE FROM {database}.{table} WHERE id={data['data']['id']};"

                    # 执行插入操作
                    cursor.execute( sql)
                    conn.commit()
                    print(f"成功执行sql语句: {sql}")
            except json.JSONDecodeError as e:
                print(f"解析 JSON 消息时出错: {e}")
            except psycopg2.Error as e:
                print(f"插入数据到 PostgreSQL 时出错: {e}")
                conn.rollback()
except KeyboardInterrupt:
    print("脚本被手动中断")
finally:
    # 关闭连接
    pubsub.unsubscribe(redis_channel)
    cursor.close()
    conn.close()
    print("已关闭 Redis 和 PostgreSQL 连接")

3.2准备Dockerfile

# 选择基础镜像
FROM python:3.9

# 设置工作目录
WORKDIR /app

# 复制 Python 代码
COPY consumer.py .

# 复制依赖文件
COPY requirements.txt .

# 安装 Python 依赖
RUN pip install --no-cache-dir -r requirements.txt

# 运行 Python 消费者
CMD ["python", "consumer.py"]

3.3 准备requirements.txt文件

redis
psycopg2

3.4 构建消费者镜像

docker build -t maxwell-python-consumer .

四、一键启动,搞定!

docker-compose up -d

五、自行测试:操作mysql表看是否同步到pg

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐