Maxwell实战之mysql实时同步mysql数据到postgres
·
Maxwell实战之mysql实时同步mysql数据到postgres
前言
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
更多推荐





所有评论(0)