Java集成Debezium时出现「未找到引擎构建器实现」错误排查
问题描述
在将Debezium集成到Java代码中,尝试通过DebeziumEngine.create(Json.class)初始化引擎以JSON格式暴露事件时,抛出DebeziumException,错误信息:
No implementation of Debezium engine builder was found
相关代码
Java初始化方法
private void setupDebezium() throws DebeziumException { log.info("Setting up debezium"); //start debezium final Properties props = new Properties(); props.setProperty("name", "engine"); props.setProperty("offset.storage", "org.apache.kafka.connect.storage.FileOffsetBackingStore"); props.setProperty("offset.storage.file.filename", "/Users/abc/mytemp1/offsets.dat"); props.setProperty("offset.flush.interval.ms", "1000"); /* begin connector properties */ props.setProperty("connector.class", "io.debezium.connector.mysql.MySqlConnector"); props.setProperty("database.hostname", "localhost"); props.setProperty("database.port", "3306"); props.setProperty("database.user", "root"); props.setProperty("database.password", "root"); props.setProperty("database.dbname", "students"); props.setProperty("database.server.id", "85744"); props.setProperty("topic.prefix", "my-app-connector"); props.setProperty("schema.history.internal", "io.debezium.storage.file.history.FileSchemaHistory"); props.setProperty("schema.history.internal.file.filename", "/Users/abc/mytemp1/schemahistory.dat"); try (DebeziumEngine<ChangeEvent<String, String>> engine = DebeziumEngine.create(Json.class) .using(props) .notifying(record -> { System.out.println(record); }).build() ) { // Run the engine asynchronously ... ExecutorService executor = Executors.newSingleThreadExecutor(); executor.execute(engine); // Do something else or wait for a signal or an event Thread.sleep(10000000); } catch (IOException e) { throw new RuntimeException(e); } catch (InterruptedException e) { throw new RuntimeException(e); } }
Maven依赖
<dependency> <groupId>io.debezium</groupId> <artifactId>debezium-api</artifactId> <version>${version.debezium}</version> </dependency> <dependency> <groupId>io.debezium</groupId> <artifactId>debezium-core</artifactId> <version>${version.debezium}</version> </dependency> <dependency> <groupId>io.debezium</groupId> <artifactId>debezium-embedded</artifactId> <version>${version.debezium}</version> </dependency> <dependency> <groupId>io.debezium</groupId> <artifactId>debezium-connector-mysql</artifactId> <version>${version.debezium}</version> </dependency>
版本信息
<version.debezium>2.1.4.Final</version.debezium>
问题原因及解决方案
1. API用法错误(核心问题)
Debezium 2.x版本中,DebeziumEngine.create()方法的参数格式已变更,不再支持直接传入Json.class,需要通过ChangeEventFormat来指定事件的序列化格式。
修改代码
将引擎初始化代码替换为以下写法,并确保导入io.debezium.engine.ChangeEventFormat类:
try (DebeziumEngine<ChangeEvent<String, String>> engine = DebeziumEngine.create(ChangeEventFormat.of(Json.class)) .using(props) .notifying(record -> { System.out.println(record); }).build() ) { // 原有异步执行逻辑保持不变 }
2. MySQL连接器属性错误
代码中使用了错误的数据库名称属性database.dbname,Debezium MySQL连接器的正确属性名是database.name,需要修改:
// 替换前 props.setProperty("database.dbname", "students"); // 替换后 props.setProperty("database.name", "students");
3. 确认依赖完整性
当前依赖已包含debezium-core、debezium-embedded等必要包,版本一致无需调整,确保构建工具(Maven)已正确下载所有依赖。
内容的提问来源于stack exchange,提问作者RusJaI
相关产品推荐
相关产品推荐

