转载请注明出处:

起因是排查两段真实代码时产生的两个疑问:

  1. AlertTriggerV2ServiceImpl.sendAlarmToThirdParty()ALARM_TO_THIRD_PARTY 推告警,如果 Kafka 是三节点集群,而某个消费者只订阅了其中一个实例,会丢消息吗?
  2. 如果用三个实例对应地址的 VIP,能订阅成功、能消费到消息吗?

结论先行:问题一的答案是"不会",问题二的答案是"能订阅成功,但不一定能消费到"。而这两个问题背后,真正的丢消息风险完全不在"订阅了几个 broker"上。本文按场景逐个拆解。


0. 先立一个正确的心智模型

很多人(包括踩过坑的我)会下意识把 Kafka 当成"按 broker 分片存储的消息队列",认为:

三个节点 = 三份不同的数据 = 要连三个才能拿全 = 连一个只能拿到 1/3

这是错的。 Kafka 的关键事实只有三条:

  1. bootstrap.servers种子列表(seed list),不是"订阅范围"。它的唯一作用是让客户端有机会开口问第一次话。
  2. 任意一个 broker 都能返回全集群视图 + 目标 topic 所有 partition 的 leader 分布在哪些 broker
  3. 数据的归属单位是 partition,不是 broker。一个 partition 的所有副本只有一份 leader 可读写,follower 只做同步复制。消费者拿到 metadata 后,直接连到每个 partition 的 leader 上拉数据。

所以:客户端"连了一个 broker"只是"用这一个地址换取了全局地图",之后的读取范围是整个 topic 的全部分区。


场景一:三节点集群,消费者只配了一个 broker 地址

现象复现

我们的配置里两种写法都存在:

# terra-consul/conf/yml/terra/lab33/terra-no,cluster.yml
spring:
  kafka:
    consumer:
      bootstrap-servers: 'kafka1:9092,kafka2:9092,kafka3:9092'   # 集群环境:列全

# terra-consul/conf/yml/terra/lab33/terra-no,dev.yml
      bootstrap-servers: 'kafka1:9092'                            # dev:只写一个

# DeviceConfigKafkaConfig.java
@Value("${spring.kafka.consumer.bootstrap-servers:kafka1:9092}")   # 兜底默认值:只写一个
private String bootstrapServers;

Go 采集端则是另一种写法:

// terra-monitor/src/cmd/monitor.go
if conf.KafkaAddr != nil {
    kafkaAddr = conf.KafkaAddr
} else {
    kafkaAddr = []string{"kafka1:9092", "kafka2:9092", "kafka3:9092"}  // 兜底列全 3 个
}
producer, err := sarama.NewSyncProducer(kafkaAddr, kafkaConfig)

结论:不会丢消息

只配 kafka1:9092,消费者依然能读到 kafka2kafka3 上作为 leader 的那些 partition。因为 metadata 里写着"partition-5 的 leader 在 172.21.254.68:9092",客户端就会自己去建那条连接。

但单实例有两个真实隐患(都不是丢消息)

① 可用性问题:重启时踩不到种子就起不来。

已经跑起来的客户端,metadata 刷新可以走任意已连接的 broker,kafka1 挂了它照样能感知 kafka2/kafka3 上的 leader 变更、继续消费。但进程重启时,如果唯一写死的种子 broker 正好不可用,客户端连开口问话的机会都没有 → 启动失败。

所以生产环境应该列全 3 个:客户端会依次尝试种子,直到有一个能用为止。

② 域名解析问题:这个现象最像"丢消息",但其实是网络。

Kafka 集群用 Docker --network host 部署,broker 之间用主机名互访:

# sdn-deploy/services/terra-kafka/deploy-cluster/terra-kafka.service_TEMPLATE.txt
-e KAFKA_ADVERTISED_LISTENERS=__KAFKA_ADVERTISED_LISTENERS__ \
-e KAFKA_INTER_BROKER_LISTENER_NAME=PLAINTEXT \

一旦消费者容器里 /etc/hosts 只配了 kafka1、没配 kafka2/kafka3,客户端会拿到 kafka2 的 leader 信息却解析不了、连不上,表现为一部分 partition 永远拉不到数据。症状是"丢了一半消息",根因是 advertised.listeners 不可达。这类问题查复制因子、查 acks 都是浪费时间,先查连通性。


场景二:用 VIP 做 bootstrap,能订阅成功吗?能消费到消息吗?

核心:Kafka 是两阶段连接

阶段 1 · bootstrap
  客户端 → bootstrap.servers 中任一地址(VIP / LB 在这里完全可用)
        → MetadataRequest → 拿到全集群 broker 列表 + partition leader 分布

阶段 2 · 数据面
  Fetch / JoinGroup / SyncGroup / FindCoordinator / OffsetCommit / Produce
        → 全部按 metadata 返回的 advertised.host:port 【另建新连接】
        → 这一步默认不经过 VIP(除非 advertised.listeners 本身就是 VIP)

由此得到唯一判据:

VIP 只影响"能否问到地图";advertised.listeners 决定"能否走到目的地"。

决策表

advertised.listeners 的配法 订阅(入组) 实际消费 典型现象
每 broker 独立真实地址,且客户端可解析可达 完全正常,VIP 变成"多余但无害"
每 broker 独立地址,但客户端不可达(跨网段 / 防火墙 / 域名不解析) ✅ 看起来成功 消费者组里能看到你的 member,但 lag 不降;日志反复出现 Connection to node 2 ... could not be established / Broker may not be available,不断 rebalance
三个 broker 全部写成同一个 VIP ⚠️ ⚠️ 不稳定 客户端按 brokerId 维护独立连接并校验对端身份,共用一个 host:port 时 L4 会把它转到随机后端 → broker id 与预期不符、连接反复断开重建、SASL 认证状态错乱、写入落到非预期分区

第三行是踩坑高发区,也是最符合直觉、最错的一种:想"用一个 VIP 屏蔽掉三个实例的差异",在 Kafka 上做不到。 Kafka 不是无状态 HTTP 服务,它有"每个 brokerId 一条专属连接"的强假设。

正确的负载均衡用法是 per-broker endpoint:3 个 VIP(或 1 个 VIP 上的 3 个不同端口),一一对应 3 个 broker,各自写进自己的 advertised.listeners

VIP 本身的实现方式也决定可用性

方式 能否用于 Kafka
纯 VRRP 漂移 VIP(keepalived 只配 virtual_ipaddress 只能当 bootstrap 用。它不是负载均衡,任一时刻只指向一台;若把 advertised 也写成它,所有流量只进一台,另外两台空转。主备切换时长连接全断:producer 靠 retries 兜住,consumer 靠 rebalance 兜住,已 ack 的数据不会因此丢
L4:LVS DR/NAT/FULLNAT、HAProxy mode tcp、nginx stream、硬件 LB 的 TCP profile 可以,但必须是连接级(connection-level) 转发。Kafka 在单条 TCP 上跑多路复用请求 + 长轮询 fetch,任何"请求级轮转 / 连接复用"都会立刻崩。DR 模式还要求客户端与 VIP 二层可达,跨网段要 FULLNAT/DNAT
L7:HTTP 代理、nginx http 块、API Gateway 完全不可用。Kafka 是私有的二进制 TCP 协议

顺带一句:我们仓库里的 keepalived 就是纯 VRRP(只配 virtual_ipaddress,没有 virtual_server/real_server 的 LVS 段),而且 conf_examples/keepalived_03_VIP1-VIP2-VIP3_V4.conf 里三个 VIP 被放进同一个 vrrp_sync_group

vrrp_sync_group TERRA_VRRP_GROUP {
    group {
        VRRP_I1     # 10.218.41.206
        VRRP_I2     # 11.83.255.214
        VRRP_I3     # 60.83.255.229
    }
}

sync group 的语义是"要么全在、要么全走",所以这三个 VIP 永远落在同一台机器上。拿它当 Kafka 的"三个入口"其实是假冗余——入口还是只有一个物理节点。

那我们该不该上 VIP?

先看现在的实际配法:

# terra-kafka.service_generator.sh
KAFKA_ADVERTISED_LISTENERS="PLAINTEXT://${TARGET_HOST_IP}:9092"   # 每节点 advertise 自己的真实 IP

每个 broker 广播的是自己的真实 IP,kafka1/2/3 只是 hosts 里的映射名。也就是说我们天然就是决策表第一行——不需要 VIP。而且 Kafka 客户端自带"轮询种子 + 优先复用已连接 broker 刷新 metadata"的容错,多写三个地址本身就是高可用,引入 VIP 只会新增一个故障域(VIP 漂移本身、LB 会话表、conntrack 超时把长连接掐掉)。

只有一种情况值得上 VIP:跨网段 / 防火墙只放行 VIP。 这时必须配套改成双 listener,并且客户端入口按 broker 区分:

listeners=PLAINTEXT_INTERNAL://<真实IP>:9093,PLAINTEXT_CLIENT://<真实IP>:9092
advertised.listeners=PLAINTEXT_INTERNAL://<真实IP>:9093,PLAINTEXT_CLIENT://<该broker专属VIP>:9092
listener.security.protocol.map=PLAINTEXT_INTERNAL:PLAINTEXT,PLAINTEXT_CLIENT:PLAINTEXT
inter.broker.listener.name=PLAINTEXT_INTERNAL

即 broker 之间走真实地址,客户端走各自独立的 VIP。


场景三:那到底什么才会丢消息?(这才是重点)

把上面两个问题排掉之后,用 kafka-topics --describe 和读代码的方式逐层看,真正能丢消息的是下面 5 处。它们和"订阅了几个 broker"一点关系都没有。

3.1 保留期只有 1 小时 —— 最容易踩的一个

# terra-kafka.service_TEMPLATE.txt
-e KAFKA_LOG_RETENTION_HOURS=1 \

消息只保留 1 小时。消费者停摆、发布回滚、或单纯落后超过 1 小时,未消费的消息已被物理清理,再配上下面这条:

props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest");

latest 意味着不会回溯历史 → 停机期间的消息 100% 真丢

仓库里其实已经有一份写对了的笔记(如何设置单个topic的消息保留策略.txt):全局 LOG_RETENTION_HOURS=72 兜底,只对高时效 topic 单独设 retention.ms=3600000。但模板里的全局值还是 1,属于"文档比配置正确"的典型。

排查建议:先用 kafka-configs --describe --entity-type topics --entity-name 实际生效的 retention,别信模板 —— 线上很可能被环境变量或 docker run 参数覆盖过。

3.2 副本因子与 ack 语义不匹配

-e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=2 \
-e KAFKA_DEFAULT_REPLICATION_FACTOR=2 \
# 没有 KAFKA_MIN_INSYNC_REPLICAS → 默认 1
位置 当前值 后果
device_config_kafka.go RequiredAcks = WaitForLocal(等价 acks=1) leader 落盘即返回成功,leader 磁盘损坏 = 数据消失
monitor.go / kafka_worker.go RequiredAcks = WaitForAll acks=all 但 min.insync.replicas=1 → ISR 可能只剩 leader,语义退化成单副本
Java spring.kafka.producer 未显式配置 kafka-clients 3.0.1 起默认幂等 + acks=all,同样被 min.insync.replicas=1 卡住
sdn-collector monitor.go 逻辑 同上

关键点:acks=all + min.insync.replicas=1acks=1。三节点集群照样能丢数据。真正的组合是 RF=3 + min.insync.replicas=2 + acks=all + enable.idempotence=true

还有一个隐蔽情况:topic 如果是在 DEFAULT_REPLICATION_FACTOR 生效之前被 auto-create 出来的(或早期单节点时代建的),Replicas 可能仍是 [1]。必须逐个核对:

kafka-topics --describe --topic device-config-update --bootstrap-server kafka1:9092
# 看每一行 Partition 的 Replicas / ISR;ISR 长度 = 1 就是单点

3.3 生产者把失败静默吞掉(本文开头那段代码)

            String alarmMsg = objectMapper.writeValueAsString(results);
            kafkaTemplate.send("ALARM_TO_THIRD_PARTY", alarmMsg);   // ① 返回的 ListenableFuture 被丢弃
        }
    } catch (Exception e) {
        log.error("kafka send alarm msg failed,  error msg is {} ", e.getMessage());   // ② 只 log,无重试无落库
    }

两个问题:

  • send()异步的。broker 拒绝、leader 不可用、消息超 max.message.bytes、producer buffer 满超时……这些失败只体现在被丢弃的那个 Future 里,外层 catch 抓不到
  • ② 告警是一次性消息,不像设备配置有 SNMP 周期轮询自愈。这条路径是真丢且不留痕迹

顺带一个 bug:InfluxDB 查询和 send 共用同一个 try。查询抛异常时 results 为 null,代码仍会往下 writeValueAsString(null),给第三方发出一条字面量 "null" 的脏消息。

正确写法:

kafkaTemplate.send(TOPIC, key, alarmMsg).addCallback(
        ok -> log.debug("alarm sent, offset={}", ok.getRecordOffset()),
        err -> { log.error("alarm send failed, persist for retry", err); retryService.save(params, alarmMsg); });

(另外记忆里有一条经验值得复述:消息过大时,生产者 max.request.size 和 broker 端 max.message.bytes 必须同时调,只调一边会出现"生产者通过、broker 拒绝"的诡异现象。)

3.4 消费端无条件 ack(at-most-once)

} catch (Exception e) {
    log.error("failed to process device config for host={}: {}", host, e.getMessage(), e);
} finally {
    // ★ 无论成功失败都 ack,防止单条毒消息阻塞整个消费线程
    ack.acknowledge();
}

enable.auto.commit=false + AckMode.MANUAL_IMMEDIATE 这个搭配本身是对的(精确一次提交时机),但 finally无条件 ack 把语义降级成了 at-most-once:DB 写入失败的消息永久跳过。

device-config-update 这是可接受的取舍——下一轮 SNMP 轮询会重推,数据自愈;同样的写法用在告警上就是丢数据。要不要改,取决于消息本身能不能重放,而不是"最佳实践"。

3.5 auto.offset.reset=latest

latest 在这三种时机下会跳过消息:新 group 首次上线、offset 被清理(offsets.retention.minutes 默认 7 天,消费者停超过 7 天)、topic 被重建。它是"从今往后"的语义,不是"从头开始"。


场景四:一个高频误判 —— "我的实例没收到这条消息"≠ 丢消息

@KafkaListener(topics = "device-config-update",
               groupId = "terra-monitor-server-device-config", ...)
// DeviceConfigKafkaConfig
factory.setConcurrency(6);   // 6 个消费者线程,匹配 9 分区

同一个 group.id 下,partition 会在实例间分摊,每条消息只被组内一个实例消费。所以:

  • terra-monitor-server 部署多副本时,单个实例日志里只看到一部分设备 → 正常,不是丢。
  • 如果 topic 分区数 < 消费者线程数,多余线程空转 → 不丢,只是浪费。
  • 如果你真的想要"每个实例都收到全量"(广播语义),必须用不同的 group.id。只连不同 broker 是没用的——这恰好又回到"broker 不是订阅单位,partition + group 才是"。

排查时对应的命令:

kafka-consumer-groups --describe --group terra-monitor-server-device-config --bootstrap-server kafka1:9092
# 看每个 partition 被哪个 CONSUMER-ID(哪个实例)持有,以及 LAG

一页速查表

问题 答案
只订阅一个 broker 会丢消息吗 不会。bootstrap 只是种子,metadata 会引导客户端读到全部分区
那要不要列全 3 个 。避免种子 broker 不可用时进程重启失败;同时能规避"advertised 主机名解析不全"的一部分表现
用 VIP 能订阅成功吗 (MetadataRequest 发给任意 broker 都返回全量)
用 VIP 能消费到消息吗 取决于 advertised.listeners 是否客户端可达。别拿"订阅成功"当连通性判据
三个 broker 共用一个 VIP 行吗 不行。客户端按 brokerId 建独立连接,共用地址会导致 id 不匹配、反复断连
什么 VIP 能用 L4(LVS/HAProxy tcp/nginx stream),且连接级转发;L7 一律不行
真正的丢消息风险 保留期 1h + latestacks=1 / min.insync.replicas=1;异步 send 无回调;无条件 ack;单种子 broker + advertised 不可达
Kafka 需要 LB 提升可用性吗 不需要,客户端内建种子轮询与 metadata 刷新。仅在跨网段/防火墙场景才考虑,且必须 per-broker endpoint

上手先跑这三条命令

# 1) 副本与 ISR:ISR 长度 = 1 就是单点
kafka-topics --describe --topic ALARM_TO_THIRD_PARTY --bootstrap-server kafka1:9092

# 2) 实际生效的 retention 与 broker 默认值(别信模板,信线上)
kafka-configs --describe --entity-type topics --entity-name ALARM_TO_THIRD_PARTY --bootstrap-server kafka1:9092
kafka-configs --describe --entity-type brokers --entity-default --bootstrap-server kafka1:9092

# 3) 消费组落点与 lag:判断"没收到"是分摊还是真丢
kafka-consumer-groups --describe --group terra-monitor-server-device-config --bootstrap-server kafka1:9092

修复清单(按优先级)

优先级 动作
P0 KAFKA_LOG_RETENTION_HOURS 调到 72 兜底,高时效 topic 单独设 retention.ms
P0 告警 send() 加回调 + 失败落库重试;把 InfluxDB 查询与 Kafka 发送的 try 分离,避免发出 "null"
P1 KAFKA_DEFAULT_REPLICATION_FACTOR=3 + KAFKA_MIN_INSYNC_REPLICAS=2;Java producer 显式 acks: all + enable-idempotence: true;Go 侧 WaitForLocalWaitForAll
P1 生产环境 bootstrap 列全 3 个地址(含 DeviceConfigKafkaConfig@Value 默认值、consul 里的 *,dev.yml);确认 kafka1/2/3 在所有消费者容器内可解析
P2 已存在但 RF=1 的 topic 用 kafka-reassign-partitions 补副本
P2 消费失败的消息进 DLQ 而非盲 ack;给告警 topic 指定 key 保证同设备有序

总结

三个可以带走的结论:

  1. broker 不是订阅单位,partition + consumer group 才是。 "只连一个 broker 怕丢消息"和"多连几个 broker 就能多收消息"是同一个误解的两面。搞清 metadata 引导机制,这类焦虑可以直接消除。

  2. VIP 只存在于连接建立的第一瞬间,之后就没它事了。 所以"VIP 配好了却收不到消息"几乎必然是 advertised.listeners 可达性问题,而不是配置写错了 bootstrap。反过来,试图用 VIP 把三个 broker 收敛成一个入口,是在跟 Kafka 的连接模型对着干——正确姿势是 per-broker endpoint。

  3. 丢消息发生在三个你不太会注意的地方:保留期、副本与 ack 的组合、以及被 catch (Exception) 吞掉的异步失败。 这三个都跟集群规模无关——3 节点集群 + retention=1h + min.insync.replicas=1 + 无回调 send(),比单节点但配置自洽的系统更容易丢数据。集群数量给人安全感,参数才决定数据命运。

判断顺序也顺便记一下:先查连通性(advertised 可达)→ 再查保留期与 offset(还在不在)→ 再查副本与 ack(有没有落对地方)→ 最后查代码里的异常处理和 ack 时机(有没有被静默跳过)。


原文地址: http://www.cveoy.top/t/topic/qHzF 著作权归作者所有。请勿转载和采集!

免费AI点我,无需注册和登录