本文目录导读:

#!/bin/bash
# Debezium 自动配置脚本
# 支持 MySQL 和 PostgreSQL 数据源
# 适用于 Debezium 1.9+ / 2.x + Kafka Connect
set -euo pipefail
# ---------- 颜色定义 ----------
RED='\033[0;31m'
GREEN='\033[0;32m'
YELLOW='\033[1;33m'
CYAN='\033[0;36m'
NC='\033[0m' # No Color
# ---------- 默认变量 ----------
CONFIG_FILE=""
MODE="interactive" # interactive / silent
SOURCE_TYPE=""
CONNECTOR_NAME="debezium-connector"
KAFKA_CONNECT_URL="http://localhost:8083"
KAFKA_BOOTSTRAP_SERVERS="localhost:9092"
DATABASE_HOSTNAME="localhost"
DATABASE_PORT=""
DATABASE_USER=""
DATABASE_PASSWORD=""
DATABASE_NAME=""
DATABASE_SERVER_ID="1" # MySQL 专用
TABLE_INCLUDE_LIST=""
TABLE_EXCLUDE_LIST=""
SNAPSHOT_MODE="initial"
OFFSET_STORAGE_TOPIC="connect-offsets"
SCHEMA_CHANGE_TOPIC="schema-changes"
SSL_ENABLED="false"
SSL_KEYSTORE_PATH=""
SSL_KEYSTORE_PASSWORD=""
SSL_TRUSTSTORE_PATH=""
SSL_TRUSTSTORE_PASSWORD=""
TRANSFORMS_ENABLED="false"
TRANSFORM_EXTRACT_NEW_RECORD="false"
PLUGIN_PATH="/kafka/connect"
# ---------- 帮助信息 ----------
usage() {
echo -e "${CYAN}Debezium 自动配置脚本${NC}"
echo ""
echo "用法: $0 [选项]"
echo ""
echo "选项:"
echo " -h, --help 显示此帮助信息"
echo " -m, --mode <mode> 运行模式:interactive(交互式,默认)/ silent(静默)"
echo " -c, --config <file> 静默模式下的配置文件路径"
echo " --source-type <type> 数据源类型:mysql / postgresql"
echo " --connector-name <name> 连接器名称(默认:debezium-connector)"
echo " --kafka-connect <url> Kafka Connect REST API 地址(默认:http://localhost:8083)"
echo " --bootstrap-servers <s> Kafka 引导服务器地址(默认:localhost:9092)"
echo " --database-host <host> 数据库主机地址(默认:localhost)"
echo " --database-port <port> 数据库端口"
echo " --database-user <user> 数据库用户"
echo " --database-password <password> 数据库密码"
echo " --database-name <db> 数据库名称"
echo " --mysql-server-id <id> MySQL 服务器 ID(默认:1)"
echo " --table-include <list> 包含的表列表(逗号分隔)"
echo " --table-exclude <list> 排除的表列表(逗号分隔)"
echo " --snapshot-mode <mode> 快照模式:initial / never / schema_only(默认:initial)"
echo " --ssl-enabled 启用 SSL"
echo " --ssl-keystore <path> SSL 密钥库路径"
echo " --ssl-truststore <path> SSL 信任库路径"
echo " --transforms-enabled 启用数据转换"
echo " --transform-new-record 启用提取新记录状态转换"
echo " --plugin-path <path> Debezium 插件路径(默认:/kafka/connect)"
echo ""
echo "示例:"
echo " 交互模式:"
echo " sudo $0 -m interactive"
echo ""
echo " 静默模式(使用配置文件):"
echo " sudo $0 -m silent -c /path/to/debezium-config.conf"
echo ""
echo " 静默模式(命令行参数):"
echo " sudo $0 --source-type mysql \\"
echo " --database-host 192.168.1.100 \\"
echo " --database-port 3306 \\"
echo " --database-user debezium \\"
echo " --database-password secret \\"
echo " --database-name mydb \\"
echo " --table-include \"public.users,public.orders\""
exit 0
}
# ---------- 日志函数 ----------
log_info() { echo -e "${GREEN}[INFO]${NC} $1"; }
log_warn() { echo -e "${YELLOW}[WARN]${NC} $1"; }
log_error() { echo -e "${RED}[ERROR]${NC} $1"; }
log_step() { echo -e "${CYAN}[STEP]${NC} $1"; }
# ---------- 参数解析 ----------
parse_args() {
while [[ $# -gt 0 ]]; do
case "$1" in
-h|--help) usage ;;
-m|--mode) MODE="$2"; shift 2 ;;
-c|--config) CONFIG_FILE="$2"; shift 2 ;;
--source-type) SOURCE_TYPE="$2"; shift 2 ;;
--connector-name) CONNECTOR_NAME="$2"; shift 2 ;;
--kafka-connect) KAFKA_CONNECT_URL="$2"; shift 2 ;;
--bootstrap-servers) KAFKA_BOOTSTRAP_SERVERS="$2"; shift 2 ;;
--database-host) DATABASE_HOSTNAME="$2"; shift 2 ;;
--database-port) DATABASE_PORT="$2"; shift 2 ;;
--database-user) DATABASE_USER="$2"; shift 2 ;;
--database-password) DATABASE_PASSWORD="$2"; shift 2 ;;
--database-name) DATABASE_NAME="$2"; shift 2 ;;
--mysql-server-id) DATABASE_SERVER_ID="$2"; shift 2 ;;
--table-include) TABLE_INCLUDE_LIST="$2"; shift 2 ;;
--table-exclude) TABLE_EXCLUDE_LIST="$2"; shift 2 ;;
--snapshot-mode) SNAPSHOT_MODE="$2"; shift 2 ;;
--ssl-enabled) SSL_ENABLED="true"; shift ;;
--ssl-keystore) SSL_KEYSTORE_PATH="$2"; shift 2 ;;
--ssl-truststore) SSL_TRUSTSTORE_PATH="$2"; shift 2 ;;
--transforms-enabled) TRANSFORMS_ENABLED="true"; shift ;;
--transform-new-record) TRANSFORM_EXTRACT_NEW_RECORD="true"; shift ;;
--plugin-path) PLUGIN_PATH="$2"; shift 2 ;;
*) log_error "未知参数: $1"; usage ;;
esac
done
}
# ---------- 加载配置文件 ----------
load_config() {
if [[ -f "$CONFIG_FILE" ]]; then
log_info "加载配置文件: $CONFIG_FILE"
source "$CONFIG_FILE"
elif [[ "$MODE" == "silent" ]]; then
log_error "静默模式需要配置文件,但文件不存在: $CONFIG_FILE"
exit 1
fi
}
# ---------- 交互式输入 ----------
interactive_input() {
echo ""
echo "========== Debezium 交互式配置 =========="
echo ""
# 数据源类型
echo -e "${CYAN}选择数据源类型:${NC}"
echo "1) MySQL"
echo "2) PostgreSQL"
read -p "请输入选项 (1 或 2): " SOURCE_CHOICE
case "$SOURCE_CHOICE" in
1) SOURCE_TYPE="mysql" ;;
2) SOURCE_TYPE="postgresql" ;;
*) log_error "无效选择"; exit 1 ;;
esac
# Kafka Connect 地址
read -p "Kafka Connect REST API 地址 [${KAFKA_CONNECT_URL}]: " input
KAFKA_CONNECT_URL="${input:-$KAFKA_CONNECT_URL}"
# 数据库连接信息
read -p "数据库主机地址 [${DATABASE_HOSTNAME}]: " input
DATABASE_HOSTNAME="${input:-$DATABASE_HOSTNAME}"
if [[ -z "$DATABASE_PORT" ]]; then
[[ "$SOURCE_TYPE" == "mysql" ]] && DATABASE_PORT="3306" || DATABASE_PORT="5432"
fi
read -p "数据库端口 [${DATABASE_PORT}]: " input
DATABASE_PORT="${input:-$DATABASE_PORT}"
read -p "数据库用户名: " DATABASE_USER
read -s -p "数据库密码: " DATABASE_PASSWORD
echo ""
read -p "数据库名称: " DATABASE_NAME
# MySQL 专用配置
if [[ "$SOURCE_TYPE" == "mysql" ]]; then
read -p "MySQL 服务器 ID [${DATABASE_SERVER_ID}]: " input
DATABASE_SERVER_ID="${input:-$DATABASE_SERVER_ID}"
fi
# 表过滤
read -p "包含的表(逗号分隔,留空表示全部): " TABLE_INCLUDE_LIST
read -p "排除的表(逗号分隔,留空表示无): " TABLE_EXCLUDE_LIST
# 快照模式
echo -e "${CYAN}选择快照模式:${NC}"
echo "1) initial (首次启动时拍摄快照)"
echo "2) never (从不拍摄快照)"
echo "3) schema_only (仅捕获 schema)"
read -p "请输入选项 (1-3) [1]: " SNAPSHOT_CHOICE
case "${SNAPSHOT_CHOICE:-1}" in
1) SNAPSHOT_MODE="initial" ;;
2) SNAPSHOT_MODE="never" ;;
3) SNAPSHOT_MODE="schema_only" ;;
*) SNAPSHOT_MODE="initial" ;;
esac
# SSL 配置
read -p "启用 SSL 连接?(y/N): " SSL_INPUT
if [[ "$SSL_INPUT" =~ ^[Yy]$ ]]; then
SSL_ENABLED="true"
read -p "SSL 密钥库路径: " SSL_KEYSTORE_PATH
read -s -p "SSL 密钥库密码: " SSL_KEYSTORE_PASSWORD
echo ""
read -p "SSL 信任库路径: " SSL_TRUSTSTORE_PATH
read -s -p "SSL 信任库密码: " SSL_TRUSTSTORE_PASSWORD
echo ""
fi
# 数据转换
read -p "启用数据转换 (Transforms)?(y/N): " TRANSFORM_INPUT
if [[ "$TRANSFORM_INPUT" =~ ^[Yy]$ ]]; then
TRANSFORMS_ENABLED="true"
read -p "启用提取新记录状态 (ExtractNewRecordState)?(y/N): " EXTRACT_INPUT
[[ "$EXTRACT_INPUT" =~ ^[Yy]$ ]] && TRANSFORM_EXTRACT_NEW_RECORD="true"
fi
echo ""
echo "========== 配置完成 =========="
echo ""
}
# ---------- 验证配置 ----------
validate_config() {
local errors=0
if [[ -z "$SOURCE_TYPE" ]] || { [[ "$SOURCE_TYPE" != "mysql" ]] && [[ "$SOURCE_TYPE" != "postgresql" ]]; }; then
log_error "数据源类型必须为 mysql 或 postgresql"
errors=$((errors+1))
fi
if [[ -z "$KAFKA_CONNECT_URL" ]]; then
log_error "Kafka Connect URL 不能为空"
errors=$((errors+1))
fi
if [[ -z "$DATABASE_HOSTNAME" ]] || [[ -z "$DATABASE_PORT" ]] || \
[[ -z "$DATABASE_USER" ]] || [[ -z "$DATABASE_PASSWORD" ]] || [[ -z "$DATABASE_NAME" ]]; then
log_error "数据库连接信息不完整"
errors=$((errors+1))
fi
if [[ "$SOURCE_TYPE" == "mysql" ]] && [[ -z "$DATABASE_SERVER_ID" ]]; then
log_error "MySQL 需要配置 Server ID"
errors=$((errors+1))
fi
if [[ "$SSL_ENABLED" == "true" ]]; then
if [[ -z "$SSL_KEYSTORE_PATH" ]] || [[ -z "$SSL_KEYSTORE_PASSWORD" ]]; then
log_error "SSL 已启用但密钥库配置不完整"
errors=$((errors+1))
fi
fi
if [[ $errors -gt 0 ]]; then
log_error "配置文件存在 $errors 个错误,请修正后重试"
exit 1
fi
log_info "配置验证通过"
}
# ---------- 生成 JSON 配置 ----------
generate_config() {
log_step "生成 Debezium 配置..."
# 基础配置
local config="{
\"name\": \"${CONNECTOR_NAME}\",
\"config\": {
\"connector.class\": \"io.debezium.connector.${SOURCE_TYPE}.${SOURCE_TYPE^}Connector\",
\"tasks.max\": \"1\",
\"database.hostname\": \"${DATABASE_HOSTNAME}\",
\"database.port\": \"${DATABASE_PORT}\",
\"database.user\": \"${DATABASE_USER}\",
\"database.password\": \"${DATABASE_PASSWORD}\",
\"database.dbname\": \"${DATABASE_NAME}\",
\"database.server.name\": \"${CONNECTOR_NAME}\",
\"topic.prefix\": \"${CONNECTOR_NAME}\",
\"snapshot.mode\": \"${SNAPSHOT_MODE}\",
\"offset.storage.topic\": \"${OFFSET_STORAGE_TOPIC}\",
\"schema.history.internal.kafka.bootstrap.servers\": \"${KAFKA_BOOTSTRAP_SERVERS}\",
\"schema.history.internal.kafka.topic\": \"${SCHEMA_CHANGE_TOPIC}\",
\"plugin.name\": \"pgoutput\",
\"plugin.path\": \"${PLUGIN_PATH}\"
"
# MySQL 额外配置
if [[ "$SOURCE_TYPE" == "mysql" ]]; then
config+=",
\"database.server.id\": \"${DATABASE_SERVER_ID}\",
\"include.schema.changes\": \"true\",
\"database.include.list\": \"${DATABASE_NAME}\"
"
fi
# PostgreSQL 额外配置
if [[ "$SOURCE_TYPE" == "postgresql" ]]; then
config+=",
\"database.dbname\": \"${DATABASE_NAME}\",
\"slot.name\": \"${CONNECTOR_NAME}_slot\",
\"publication.name\": \"${CONNECTOR_NAME}_publication\",
\"publication.autocreate.mode\": \"filtered\"
"
fi
# 表过滤
if [[ -n "$TABLE_INCLUDE_LIST" ]]; then
config+=",
\"table.include.list\": \"${TABLE_INCLUDE_LIST}\"
"
fi
if [[ -n "$TABLE_EXCLUDE_LIST" ]]; then
config+=",
\"table.exclude.list\": \"${TABLE_EXCLUDE_LIST}\"
"
fi
# SSL 配置
if [[ "$SSL_ENABLED" == "true" ]]; then
config+=",
\"database.ssl.mode\": \"verify_ca\",
\"database.ssl.keystore\": \"${SSL_KEYSTORE_PATH}\",
\"database.ssl.keystore.password\": \"${SSL_KEYSTORE_PASSWORD}\"
"
if [[ -n "$SSL_TRUSTSTORE_PATH" ]]; then
config+=",
\"database.ssl.truststore\": \"${SSL_TRUSTSTORE_PATH}\",
\"database.ssl.truststore.password\": \"${SSL_TRUSTSTORE_PASSWORD}\"
"
fi
fi
# 数据转换(Transforms)
if [[ "$TRANSFORMS_ENABLED" == "true" ]]; then
config+=",
\"transforms\": \"unwrap\",
\"transforms.unwrap.type\": \"io.debezium.transforms.ExtractNewRecordState\",
\"transforms.unwrap.drop.tombstones\": \"false\",
\"transforms.unwrap.delete.handling.mode\": \"rewrite\"
"
if [[ "$TRANSFORM_EXTRACT_NEW_RECORD" == "true" ]]; then
config+=",
\"transforms.unwrap.add.source.fields\": \"table,lsn,ts_ms\"
"
fi
fi
config+="
}
}"
CONFIG_JSON="$config"
log_info "配置生成完成"
}
# ---------- 检查 Kafka Connect 健康状态 ----------
check_kafka_connect() {
log_step "检查 Kafka Connect 状态..."
local max_retries=5
local retry_count=0
while [[ $retry_count -lt $max_retries ]]; do
if curl -sf "$KAFKA_CONNECT_URL" > /dev/null 2>&1; then
log_info "Kafka Connect 服务正常"
return 0
else
retry_count=$((retry_count+1))
if [[ $retry_count -lt $max_retries ]]; then
log_warn "Kafka Connect 不可用,等待 5 秒后重试... ($retry_count/$max_retries)"
sleep 5
fi
fi
done
log_error "Kafka Connect 服务无法连接,请确保服务已启动: $KAFKA_CONNECT_URL"
exit 1
}
# ---------- 部署连接器 ----------
deploy_connector() {
log_step "部署 Debezium 连接器..."
# 检查是否已存在同名连接器
local existing
existing=$(curl -s -o /dev/null -w "%{http_code}" "${KAFKA_CONNECT_URL}/connectors/${CONNECTOR_NAME}")
if [[ "$existing" == "200" ]]; then
log_warn "连接器 '${CONNECTOR_NAME}' 已存在,将更新配置..."
local response
response=$(curl -s -X PUT \
-H "Content-Type: application/json" \
-d "$CONFIG_JSON" \
"${KAFKA_CONNECT_URL}/connectors/${CONNECTOR_NAME}/config")
if echo "$response" | grep -q "error_code"; then
log_error "更新连接器失败: $(echo $response | python3 -c "import sys,json; print(json.load(sys.stdin).get('message',''))" 2>/dev/null || echo "$response")"
exit 1
fi
log_info "连接器已更新"
else
local response
response=$(curl -s -X POST \
-H "Content-Type: application/json" \
-d "$CONFIG_JSON" \
"${KAFKA_CONNECT_URL}/connectors")
if echo "$response" | grep -q "error_code"; then
log_error "创建连接器失败: $(echo $response | python3 -c "import sys,json; print(json.load(sys.stdin).get('message',''))" 2>/dev/null || echo "$response")"
exit 1
fi
log_info "连接器已创建"
fi
}
# ---------- 验证连接器状态 ----------
verify_connector() {
log_step "验证连接器状态..."
sleep 3
local status
status=$(curl -s "${KAFKA_CONNECT_URL}/connectors/${CONNECTOR_NAME}/status")
local state
state=$(echo "$status" | python3 -c "import sys,json; print(json.load(sys.stdin).get('connector',{}).get('state','UNKNOWN'))" 2>/dev/null || echo "UNKNOWN")
case "$state" in
"RUNNING")
log_info "连接器状态: ${GREEN}${state}${NC}"
local tasks
tasks=$(echo "$status" | python3 -c "import sys,json; tasks=json.load(sys.stdin).get('tasks',[]); print('\n'.join([f'{t[\"id\"]}:{t[\"state\"]}' for t in tasks]))" 2>/dev/null || echo "N/A")
log_info "任务状态:\n$tasks"
;;
"PAUSED")
log_warn "连接器状态: ${YELLOW}${state}${NC}"
;;
"FAILED")
log_error "连接器状态: ${RED}${state}${NC}"
local trace
trace=$(echo "$status" | python3 -c "import sys,json; print(json.load(sys.stdin).get('connector',{}).get('trace',''))" 2>/dev/null)
[[ -n "$trace" ]] && log_error "错误追踪:\n$trace"
exit 1
;;
*)
log_warn "连接器状态: ${YELLOW}${state}${NC}"
;;
esac
}
# ---------- 显示连接器配置摘要 ----------
show_summary() {
echo ""
echo "========== Debezium 连接器部署总结 =========="
echo "连接器名称: ${CONNECTOR_NAME}"
echo "数据源类型: ${SOURCE_TYPE}"
echo "数据库: ${DATABASE_HOSTNAME}:${DATABASE_PORT}/${DATABASE_NAME}"
echo "Kafka Connect: ${KAFKA_CONNECT_URL}"
echo "快照模式: ${SNAPSHOT_MODE}"
echo "SSL 启用: ${SSL_ENABLED}"
echo "数据转换: ${TRANSFORMS_ENABLED}"
echo "包含表: ${TABLE_INCLUDE_LIST:-全部}"
echo "排除表: ${TABLE_EXCLUDE_LIST:-无}"
echo ""
echo "管理命令:"
echo " 查看状态: curl ${KAFKA_CONNECT_URL}/connectors/${CONNECTOR_NAME}/status"
echo " 暂停: curl -X PUT ${KAFKA_CONNECT_URL}/connectors/${CONNECTOR_NAME}/pause"
echo " 恢复: curl -X PUT ${KAFKA_CONNECT_URL}/connectors/${CONNECTOR_NAME}/resume"
echo " 删除: curl -X DELETE ${KAFKA_CONNECT_URL}/connectors/${CONNECTOR_NAME}"
echo "============================================"
}
# ---------- 主流程 ----------
main() {
echo -e "${CYAN}==================================================${NC}"
echo -e "${CYAN} Debezium 自动配置脚本 v1.0${NC}"
echo -e "${CYAN}==================================================${NC}"
# 解析参数并加载配置
parse_args "$@"
if [[ "$MODE" == "silent" ]]; then
load_config
else
interactive_input
fi
# 验证配置
validate_config
# 检查依赖
if ! command -v curl &> /dev/null; then
log_error "请先安装 curl"
exit 1
fi
# 生成配置
generate_config
# 部署流程
check_kafka_connect
deploy_connector
verify_connector
show_summary
log_info "Debezium 连接器配置完成!"
}
# ---------- 执行 ----------
main "$@"
使用说明
脚本准备
# 赋予执行权限 chmod +x debezium-auto-config.sh
交互式使用(推荐新手)
sudo ./debezium-auto-config.sh -m interactive
脚本会引导您逐步完成所有配置。
静默模式(使用配置文件)
准备配置文件 debezium-config.conf:
# 示例配置文件 SOURCE_TYPE="mysql" KAFKA_CONNECT_URL="http://localhost:8083" DATABASE_HOSTNAME="192.168.1.100" DATABASE_PORT="3306" DATABASE_USER="debezium" DATABASE_PASSWORD="your_password" DATABASE_NAME="your_database" TABLE_INCLUDE_LIST="public.users,public.orders" SNAPSHOT_MODE="initial" SSL_ENABLED="false"
执行静默安装:
sudo ./debezium-auto-config.sh -m silent -c debezium-config.conf
完全命令行参数模式
sudo ./debezium-auto-config.sh --source-type mysql \
--database-host 192.168.1.100 \
--database-port 3306 \
--database-user debezium \
--database-password Secret123 \
--database-name mydb \
--table-include "users,orders" \
--snapshot-mode initial
前置条件
- 已运行的 Kafka Connect 集群(REST API 默认端口 8083)
- Debezium 连接器插件已部署到 Kafka Connect 的
plugin.path - 数据库配置:
- MySQL:启用 binlog(
log_bin=ON)且使用ROW格式 - PostgreSQL:设置
wal_level=logical并创建复制角色
- MySQL:启用 binlog(
- 网络可达:脚本运行节点可访问 Kafka Connect 和数据库
脚本功能特性
- 支持 MySQL 和 PostgreSQL 数据源
- 自动检测并更新已存在的连接器
- 详细的错误处理和状态验证
- SSL/TLS 加密连接支持
- 表过滤(包含/排除)
- 数据转换(ExtractNewRecordState)和常用管理命令输出
错误排查
若连接器状态为 FAILED,可查看详细日志:
# 查看连接器错误 curl http://localhost:8083/connectors/debezium-connector/status | jq . # 查看 Kafka Connect 日志 docker logs kafka-connect-container # 容器部署 journalctl -u kafka-connect # 系统服务部署