You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kafka Connect Redis Sink无法正常工作问题排查求助

PostgreSQL -> Kafka Connect -> Redis 数据流转架构故障排查

架构概述

  • 数据源:PostgreSQL,使用io.debezium.connector.postgresql.PostgresConnector
  • 数据接收端(Sink):Redis,使用com.github.jcustenborder.kafka.connect.redis.RedisSinkConnector

Docker Compose 服务配置

version: "3.8"
services:
  content_sql:
    container_name: content_sql
    image: postgres:13-alpine
    environment:
      POSTGRES_USER: content_admin
      POSTGRES_PASSWORD: 123456xyz
      POSTGRES_DB: madison_content
    command:
      - "postgres"
      - "-c"
      - "wal_level=logical"
    ports:
      - 5432:5432
    restart: always

  redis:
    container_name: redis
    image: redis:7-alpine
    ports:
      - 6379:6379

  zookeeper:
    container_name: zookeeper
    image: zookeeper:3.8-temurin
    environment:
      JVMFLAGS: -XX:+UseG1GC -XX:+DisableExplicitGC
    ports:
      - 2181:2181
      - 2888:2888
      - 3888:3888
    restart: always

  kafka:
    container_name: kafka
    image: ubuntu/kafka:edge
    environment:
      TZ: UTC
      ZOOKEEPER_HOST: host.docker.internal
      ZOOKEEPER_PORT: 2181
    ports:
      - 9092:9092
    depends_on:
      - zookeeper
    restart: always

  kafka_connector:
    container_name: kafka_connector
    image: debezium/connect-base:latest
    environment:
      GROUP_ID: 1
      CONFIG_STORAGE_TOPIC: content_configs
      OFFSET_STORAGE_TOPIC: cache_content
      STATUS_STORAGE_TOPIC: connect_statuses
      BOOTSTRAP_SERVERS: kafka:9092
      KEY_CONVERTER: org.apache.kafka.connect.storage.StringConverter
      VALUE_CONVERTER: org.apache.kafka.connect.storage.StringConverter
      KAFKA_CONNECT_PLUGINS_DIR: /usr/share/local-connectors
    volumes:
      - ~/Documents/kafka-plugins:/usr/share/local-connectors
    ports:
      - 8083:8083
    depends_on:
      - content_sql
      - kafka
      - redis
      - zookeeper
    restart: always

连接器注册配置

PostgreSQL 源连接器

{
    "name": "postgres-source",
    "config": {
        "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
        "task.max": "3",
        "database.hostname": "content_sql",
        "database.port": "5432",
        "database.user": "content_admin",
        "database.password": "123456xyz",
        "database.dbname": "madison_content",
        "database.whitelist": "madison_content",
        "table.include.list": "public.thread",
        "topic.prefix": "cache",
        "schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
        "plugin.name": "pgoutput"
    }
}

Redis Sink 连接器

{
    "name": "redis-sink",
    "config": {
        "connector.class": "com.github.jcustenborder.kafka.connect.redis.RedisSinkConnector",
        "redis.hosts": "redis:6379",
        "redis.database": "1",
        "redis.client.mode": "Standalone",
        "task.max": "3",
        "key.converter": "org.apache.kafka.connect.storage.StringConverter",
        "internal.key.converter": "org.apache.kafka.connect.storage.StringConverter",
        "value.converter": "org.apache.kafka.connect.storage.StringConverter",
        "internal.value.converter": "org.apache.kafka.connect.storage.StringConverter",
        "plugin.path": "/usr/share/local-connectors",
        "topics.regex": "cache.(.*)"
    }
}

当前问题与疑问

  • 所有服务通过docker-compose启动正常,插件已正确加载。
  • 向PostgreSQL插入数据后,日志显示有记录产生,手动检查Kafka主题和消息,确认主题已创建且消息存在。
  • Redis Sink始终无法正常工作,使用redis-cli的KEYS *命令检查Redis数据是否正确?请协助排查问题。

内容的提问来源于stack exchange,提问作者ThaiLB

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.06 19:40:38