自动配置Debezium的脚本

wen 实用脚本 23

本文目录导读:

自动配置Debezium的脚本

  1. 使用说明
  2. 前置条件
  3. 脚本功能特性
  4. 错误排查
#!/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

前置条件

  1. 已运行的 Kafka Connect 集群(REST API 默认端口 8083)
  2. Debezium 连接器插件已部署到 Kafka Connect 的 plugin.path
  3. 数据库配置
    • MySQL:启用 binlog(log_bin=ON)且使用 ROW 格式
    • PostgreSQL:设置 wal_level=logical 并创建复制角色
  4. 网络可达:脚本运行节点可访问 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          # 系统服务部署

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