Kafka producer 构造失败:fat jar 与 commonPool 的类加载问题
目录
生产环境里一台实例开始刷 Failed to construct kafka producer,同一集群的其余实例一切正常。重启之后错误消失,几周后又随某次部署再次出现。这个报错我们前几次都记成「Kafka 集群 / Schema Registry 抖动」——直到这次把它拆开,才发现根因不在网络:一次发生在错误线程上的类初始化,让这个实例上一整批 topic 的发送持续失败。
背景#
- 一个对外提供 HTTP API 的 Spring Boot 服务,打成 fat jar 用
java -jar跑在容器里;同集群多实例,滚动发布是 start-first:新实例先起来、旧实例再退,所以新实例头几十秒基本没有流量; - 消息通过 spring-kafka 的
KafkaTemplate发出,而且不止一个:一部分 topic 的 value 用 Confluent 的 JSON Schema 序列化器(对接 Schema Registry),另一部分用普通的StringSerializer。这个区分后面很关键;版本组合是 Boot 3.0.x、kafka-clients 3.4.0、Confluent serializer 7.4.0; - 发送入口不止 HTTP 请求路径一处:还有一个
@Scheduled重试任务,在应用就绪几秒后开始、每隔几秒扫一次积压队列,用CompletableFuture.runAsync(...)异步补发; - 所有发送都包在
try { send } catch { log }里,失败只记日志、不打断业务,丢消息是静默的,不翻日志没人知道。
现象:单个实例全挂,其余实例正常#
- 报错集中在一个实例上,同集群的其余实例全程零错误;
- 中招实例上,走 Confluent 序列化器的 topic 全部发不出去,业务 topic、日志 topic、DLQ 一起挂;
- 重启应用后就恢复了,所以每次都被记成「重启之后 Kafka 恢复了」。
还有两个细节当时没看懂,后来成了定位的关键:
- 后续每条报错都长得一模一样,连异常信息尾部方括号里的线程名都相同,而那个线程早就不是当前报错的线程;
- 挂掉的那批 topic 一起挂,而它们共享的是同一个 Java 类,不是外部的 Schema Registry。这一点当时只是猜测,后来在复现里才把边界划清楚:同一个 JVM 里走
StringSerializer的 topic 全程正常,其中还包括一个 DLQ。也就是说坏掉的不是「这个实例的 Kafka」,而是共用那个序列化器配置类的那批 template——当时如果有人只盯着后一批看,得到的结论会是「Kafka 一切正常」。
排查:三个对不上的地方#
第一,服务端故障说不通。broker 或 Schema Registry 出问题不会精确地挑中一个实例、放过同集群的其他实例。这一条基本排除了「Kafka 抖动」。
第二,时间对不上。报错从实例启动后几秒就开始了,而那台实例当时还没接过任何 HTTP 流量。能在这个窗口里发 Kafka 的只有那个 @Scheduled 重试任务:它首跑时队列里恰好有积压,于是抢在任何请求之前,触发了这个 JVM 的第一次发送。
第三,线程名对不上。NoClassDefFoundError 的 message 里带着 [in thread "..."],那个线程和打出这条日志的线程不是同一个。查了 JVM 规范(JVMS §5.5) 确认类初始化失败后会进入 erroneous state,后续使用时抛出 NoClassDefFoundError。至于异常 message 里保留原始线程名,是我在 Temurin 25 这次复现中观察到的实现细节,不能当成 JVMS 保证的格式。
根因:四个环节怎么凑到一起的#
Confluent 的配置类在 clinit 里按名解析类#
AbstractKafkaSchemaSerDeConfig 把 context.name.strategy 的默认值声明成 String 类名:
public static final String CONTEXT_NAME_STRATEGY_DEFAULT =
NullContextNameStrategy.class.getName();
...
.define(CONTEXT_NAME_STRATEGY, Type.CLASS, CONTEXT_NAME_STRATEGY_DEFAULT, ...)
而 kafka-clients 的 ConfigDef 在 define() 时就会验证默认值:Type.CLASS 的 String 默认值要经 Utils.loadClass 解析成 Class,用的是 getContextOrKafkaClassLoader():TCCL(线程上下文类加载器,Thread Context ClassLoader)非 null 就用 TCCL,只有为 null 才回退到 Kafka 自己的类加载器。这一切发生在 KafkaJsonSchemaSerializerConfig 的静态初始化(<clinit>)里。
fat jar 的类只对 Boot 的类加载器可见#
java -jar 启动的 Spring Boot 应用,业务类和全部依赖都在 BOOT-INF/ 下,由 Boot 的 LaunchedClassLoader(3.2 之前叫 LaunchedURLClassLoader)加载;系统 AppClassLoader 只看得见 fat jar 外壳。IDE 里跑、mvn spring-boot:run、单元测试用的都是平铺 classpath,所以这个问题在本地开发环境里基本碰不到,这也是它「本地全好、只在生产炸」的原因。
commonPool 线程的 TCCL 是系统类加载器#
Java 25 的 ForkJoinPool javadoc 原文:
If no thread factory is supplied via a system property, then the common pool uses a factory that uses the system class loader as the thread context class loader.
这是 JDK 9 起的设计(防应用服务器场景的类加载器泄漏)。在本文这个运行方式下,Spring 管理的常见线程(Tomcat worker、@Scheduled 调度线程、Kafka listener 容器线程)通常带有应用 ClassLoader;具体仍取决于线程创建方式和 ThreadFactory。一个常见的例外入口,就是不传 executor 的 CompletableFuture.runAsync(...) 落进去的 commonPool。我们那个重试任务用的就是裸 runAsync。
还有一跳容易漏掉:坏 TCCL 会继续往下传染。那个 commonPool 线程里又把「写请求日志」submit() 给了一个普通固定池,而 Executors.newFixedThreadPool 的 worker 是惰性创建的——new Thread() 继承的是提交者的 TCCL。于是真正执行第一次 <clinit> 的是这个二手线程,线程名还是 JDK 默认的 pool-N-thread-M,和 ForkJoinPool 看不出任何关系。上面「线程名对不上」那条线索之所以误导了我们很久,就是因为它指向的是这个名字。
类初始化只发生一次#
排查第三条线索时已经撞见失败那一半:类进入 erroneous state 后,后续初始化会失败并抛出 NoClassDefFoundError;成功初始化也只需要一次,哪个线程先完成,就决定这个 ClassLoader 之后的状态。
四个环节连起来:
(启动后几秒首跑,早于任何请求)"] --> B["CompletableFuture.runAsync(...)
未传 executor"] B --> C["ForkJoinPool.commonPool 线程
TCCL = 系统 AppClassLoader"] C --> C2["再 submit() 给一个普通固定池
惰性新建的 worker 继承坏 TCCL
线程名却是 pool-N-thread-M"] C2 --> D["本 JVM 第一次构造 KafkaProducer
触发 KafkaJsonSchemaSerializerConfig 的 clinit"] D --> E["ConfigDef.define() 验证 String 默认值
经 TCCL 解析 NullContextNameStrategy"] E --> F["系统加载器看不见 BOOT-INF/lib
ConfigException → ExceptionInInitializerError"] F --> G["类被 JVM 标记为 erroneous"] G --> H["此后所有走这个序列化器的发送:
NoClassDefFoundError
(走 StringSerializer 的 topic 不受影响)"]
复现:不需要 Kafka 集群#
这次排查用到的复现和修复,我整理成了一个可以直接跑的仓库:meirongdev/kafka-tccl-issue。除了下面这个类加载器层面的最小复现,里面还有 fat jar 上的端到端复现、平铺 classpath 的对照组(./run.sh flat)、坏 TCCL 下逐个类初始化的对照实验(./run.sh compare),以及后面三层修复各自的开关。
思路:用一个 parent 指向 platform loader 的 URLClassLoader 扮演 Boot 的类加载器(保证系统 classpath 看不见这些 jar),然后让 commonPool 线程和「TCCL 正确的线程」分别去抢第一次初始化。对比就是这两段:
// 坏线程:裸 supplyAsync 落进 commonPool,TCCL = 系统类加载器
CompletableFuture.supplyAsync(() -> init(poisoned)).join();
// 好线程:TCCL 显式设成扮演 Boot 的那个类加载器
Thread good = new Thread(() -> init(poisoned));
good.setContextClassLoader(poisoned);
good.start();
// init() 里就一句:Class.forName(CONFIG_CLASS, true, loader)
跑起来(不用起 Kafka,也不用 Spring):
git clone https://github.com/meirongdev/kafka-tccl-issue
cd kafka-tccl-issue
./run.sh minimal
在我机器上(Temurin 25.0.2,Confluent 7.4.0 + kafka-clients 3.4.0)的输出,省略了开头的环境信息和 SLF4J 的 NOP 提示:
[1] 第一次初始化发生在 commonPool 上 -> ExceptionInInitializerError
caused by ConfigException: Invalid value io.confluent.kafka.serializers.context.NullContextNameStrategy for configuration context.name.strategy: Class io.confluent.kafka.serializers.context.NullContextNameStrategy could not be found.
[2] TCCL 正确的线程,同一个 loader -> NoClassDefFoundError: Could not initialize class io.confluent.kafka.serializers.json.KafkaJsonSchemaSerializerConfig
caused by ExceptionInInitializerError: Exception org.apache.kafka.common.config.ConfigException: Invalid value io.confluent.kafka.serializers.context.NullContextNameStrategy for configuration context.name.strategy: Class io.confluent.kafka.serializers.context.NullContextNameStrategy could not be found. [in thread "ForkJoinPool.commonPool-worker-1"]
[3] 全新 loader,好线程抢先初始化 -> OK
[4] 好线程先初始化后,commonPool 再用 -> OK
四行输出对应这次排查的四个结论:
- [1] 就是生产报错的原文:
NullContextNameStrategy could not be found,类明明在 jar 里,只是 TCCL 看不见它; - [2] TCCL 正确的线程也救不回来,并且异常里重放了
[in thread "ForkJoinPool.commonPool-worker-1"]——机制和生产日志里那个「旧线程名」是同一个,只是生产上第一次失败落在固定池的 worker 上,重放出来的名字是pool-N-thread-M; - [3][4] 换一个全新的 classloader,让好线程先完成初始化,之后 commonPool 也能正常使用:先到先得。
把版本组合换成 Confluent 8.1.1 + kafka-clients 3.9.2(mvn -Dconfluent.version=8.1.1 -Dkafka.version=3.9.2 -DskipTests package 之后再跑一遍),[1] 和 [2] 的输出一字不差(8.x 的 clinit 多依赖一个 jackson-annotations,其余行为一致)。这一组我只跑了一次,没贴输出。
为什么很少有人碰到#
这个问题需要好几个条件同时成立,每一环都筛掉大部分项目:
- 只存在于 fat jar 运行形态。平铺 classpath 的本地开发、测试、
spring-boot:run都不受影响,所以开发期的验证碰不到它(./run.sh flat是这一条的反证); - 只影响用 Confluent Schema Registry 序列化器的项目。kafka-clients 自己的
ProducerConfig在 clinit 里不按名解析任何类,spring-kafka 的JsonSerializer/StringSerializer也没有这个模式(./run.sh compare是这一条的实测:同一个坏 TCCL 下逐个初始化,只有 Confluent 的配置类炸); - 需要一条落在 commonPool 上的发送路径,通常来自裸
CompletableFuture.runAsync; - 这条路径还得赢得整个 JVM 的第一次 producer 构造。大多数应用的第一次发送来自请求线程或 listener 线程,类早就被正确初始化了。我们是「启动即有活干的调度任务 + start-first 部署的无流量窗口」两件事叠出来的稳定时序;
- 中招之后重启即愈,报的人多半也就记成网络或 Schema Registry 抖动了。公开能查到的记录里,Confluent 社区论坛这个帖子的版本组合、报错线程和「间歇性、重启恢复」的描述都和我们一致(他那边是 Avro 序列化器,和我这里共用同一个父配置类
AbstractKafkaSchemaSerDeConfig),但楼里没人挖到 TCCL 这一层。
还有一个容易误判的点:这和 Spring Boot 2 还是 3 未必直接相关。本文观察到的触发条件还取决于 JDK 实现、commonPool 配置和打包形态;JDK 9+ 的默认 commonPool 使用系统类加载器作为 TCCL,所以这个问题才可能发生。JDK 8 到 JDK 9 之间的具体变更我没有逐版本去查。
怎么判断自己中不中招#
不用等线上炸。把应用打成 fat jar、用 java -jar 跑一次,在 commonPool 线程上看一眼 TCCL 能不能看见你的序列化器依赖——这个检查只有在 fat jar 形态下才有意义,IDE 和 mvn spring-boot:run 都是平铺 classpath,查不出问题:
String verdict = CompletableFuture.supplyAsync(() -> {
ClassLoader tccl = Thread.currentThread().getContextClassLoader();
try {
// initialize=false:只看可见性,不触发 clinit —— 否则这次探测本身就替好线程把类初始化了
Class.forName("io.confluent.kafka.serializers.context.NullContextNameStrategy", false, tccl);
return "看得见,这条链在你这里不成立";
} catch (ClassNotFoundException | NoClassDefFoundError e) {
return "看不见 ← 触发条件成立";
}
}).join();
initialize=false 这个参数不能省:用 true 去探测,等于替好线程完成了初始化,顺手把 bug 修掉,然后你会得到「查不出问题」的结论。
另外两件事顺手可以查:grep 一遍只传一个参数的 CompletableFuture.runAsync( / supplyAsync((以及 thenApplyAsync 这类不带 executor 的重载),它们全都落在 commonPool 上;再看看这个 JVM 的第一次 producer 构造由谁触发——启动即有活干的 @Scheduled、ApplicationRunner、消息驱动的补偿任务,都可能抢在请求线程前面。仓库里 Tccl.java、ExecutorTcclProbe.java 把这几项做成了启动横幅,./run.sh repro 一跑就能看到每类线程的 TCCL 各是什么。
升级救得了吗#
我写这篇时(2026-08-26)检查过:
- schema-registry 的 master 分支上,这些默认值仍是 String 类名,而
Type.CLASS型配置从 7.4.0 的 3 个增加到了 7 个; - kafka trunk(4.x 线)的
ConfigDef.ConfigKey仍在 define 时解析默认值,Utils.loadClass仍是 TCCL 优先、仅在 TCCL 为 null 时回退,commonPool 线程的 TCCL 非 null 但看不见 jar,回退永远不会触发; - Java 25 的
ForkJoinPooljavadoc 仍写明 common pool 用系统类加载器当 TCCL,上面的复现就是在 JDK 25 上跑出来的。
所以单纯升级 Boot、JDK 或 Confluent 版本未必能解决问题:至少我检查的版本仍保留这条触发链。对本文描述的 nested-jar 可见性问题,Spring Boot 3.3+ 的 jarmode=tools extract 默认布局把依赖放回 lib/ 并由 manifest 引用(顺便也是 CDS / AOT cache 友好的布局),可以消除这类可见性差异,仓库里 ./run.sh extract 是这一组 A/B 的实测。但下一个人可以把打包形态改回去,所以我不拿它替代代码修复。
解决方案#
真正修掉问题的是第一条,后面两条是防复发和兜底。
启动时在 main 线程抢先初始化#
利用「初始化成功也只需要一次」的语义,在任何池线程存在之前把类初始化掉:
@Configuration
public class KafkaProducerConfig {
@PostConstruct
void initSerializerConfigOnMain() throws ClassNotFoundException {
// 若这个类的第一次初始化发生在 TCCL 错误的线程上,
// 该 ClassLoader 下的类会进入 erroneous state,
// 后续 producer 构造会失败直到实例重启
Class.forName(KafkaJsonSchemaSerializerConfig.class.getName(), true,
getClass().getClassLoader());
}
}
选 @PostConstruct 而不是 ApplicationRunner 是有讲究的:@Scheduled 任务通常在 context refresh 完成后开始调度,而 @Configuration 单例的 @PostConstruct 在 refresh 期间就跑完了;runner 在 refresh 之后执行,可能与已经启动的调度器并发——把 initialDelay 调成 0 跑 ./run.sh race,就能看到第一次调度触发排在 runner 前面;这是竞态,不是必然失败。
更彻底的变体是对每个 ProducerFactory 调一次 createProducer() 做 warm-up(DefaultKafkaProducerFactory 在非 transactional、未启用 producer-per-consumer-partition 等特殊模式时,每个 factory 共享一个 producer,后续发送可以复用它),让整条序列化链上所有类的 clinit 都在 main 线程完成,代价是启动时多一次 bootstrap 地址解析,在 start-first 部署下,这反而是我们想要的 fail-fast。
给线程池固定 TCCL#
给自建线程池配一个固定 TCCL 的 ThreadFactory,顺便把线程名从 pool-9-thread-2 换成能看懂的:
private static final ClassLoader APP_CL = HttpService.class.getClassLoader();
private static final AtomicInteger SEQ = new AtomicInteger();
static final ExecutorService POOL = Executors.newFixedThreadPool(SIZE, r -> {
Thread t = new Thread(r, "outbound-log-" + SEQ.incrementAndGet());
t.setContextClassLoader(APP_CL);
return t;
});
同时治理裸 runAsync/supplyAsync,给它们显式 executor。一个容易漏的点:Spring 的 ThreadPoolTaskExecutor 默认不会替你设置固定 TCCL,worker 创建时通常继承创建线程的 TCCL;如果新加的 executor 不配上面的 factory,就是换个地方复发。
让编排器换掉坏实例#
这次真正把损失放大的是「静默」:KafkaException 是 RuntimeException,被 catch (Exception e) { log.error(...); } 接住之后业务无感、监控无感,坏实例可以带病跑很久。检测到序列化路径上的 NoClassDefFoundError / ExceptionInInitializerError 时把 liveness 打成 DOWN,就能把它变成「由编排器自动替换实例」:
// context 是注入的 ApplicationContext
AvailabilityChangeEvent.publish(context, LivenessState.BROKEN);
这行是我在 Boot 3.0.5 上实测跑通的,发布之后 /actuator/health/liveness 返回 HTTP 503。容器探针要指向 liveness 分组(management.endpoint.health.probes.enabled=true),并显式配置 health group、只包含 livenessState 和这个故障的专用 indicator;不要探测聚合的 /actuator/health,它可能带着 DB、Redis 这些外部依赖的 indicator,依赖抖一下就可能把多个实例一起重启。完整的 health group 配置和 indicator 在仓库的第三层修复里。配合 server.shutdown=graceful 把在途请求 drain 干净,比在 catch 块里到处写 System.exit(1) 更容易控制。
教训#
类初始化失败会留在 ClassLoader 上,不随下一次调用自愈:第一次初始化发生在哪个线程,决定了这个实例之后还能不能发消息。
这个形状是可以复用的:一次性的全局初始化 + 隐式的线程上下文 + 不确定的首次执行者。三者凑齐,故障就只在某些实例、某些启动顺序下出现,本地怎么试都是好的。认出形状之后,同族的地方也好找了——ServiceLoader.load(Class) 的 javadoc 写的就是 “using the current thread’s context class loader”,等价于 ServiceLoader.load(service, Thread.currentThread().getContextClassLoader()),各种 SPI 发现、provider 查找都踩在同一个假设上。
总结#
- 在这次 Temurin 25 复现里,
NoClassDefFoundErrormessage 保留了第一次失败线程名;不要把这个实现细节当成所有 JVM 的通用保证; - 「重启修好了」之后我现在会追问重启的是哪一层:我们把功记在重启 Kafka 集群上很久,所以它一直复发;
- 本地复现不了的问题,我现在先怀疑打包和运行形态的差异:平铺 classpath 和 fat jar 是两种不同的环境;
- 对默认 commonPool 而言,JDK 9+ 的实现会使用系统类加载器作为 TCCL;裸
runAsync里第一次触发的<clinit>值得多看一眼; - 类初始化只发生一次:错误线程抢先会让这个 ClassLoader 下的后续使用失败;在 main 上抢先做一次,可以避开本文描述的触发链。
相关文章#
参考资料#
- 复现与修复的仓库:meirongdev/kafka-tccl-issue
- Confluent Community — NullContextNameStrategy could not be found error intermittently
- AbstractKafkaSchemaSerDeConfig(schema-registry v7.4.0) / master 同文件
- ConfigDef.java(apache/kafka trunk) / Utils.java
- ForkJoinPool(Java 25 javadoc)
- JVMS §5.5 Initialization(erroneous state 与 NoClassDefFoundError)
- Spring Boot — Efficient deployments(jarmode tools / extract)
- Spring Boot — Kubernetes probes 与 availability state