kestra-io / kestra-io/plugin-debezium

io.kestra.plugin.debezium.mongodb.Trigger fails/unstable with MongoDB Replica Set CDC (0 messages, OpenLineage init error, MBean registration interruption)

Open
#148 0 comments 0 reactions 0 assignees View on GitHub
area/plugin good first issue
Dominant language
Java
Stars
5
Forks
11
Avg merge
2d 5h
Merged PRs (30d)
9

Description

### Describe the issue

**Issue summary**
We are using MongoDB Replica Set on Kubernetes and Kestra `io.kestra.plugin.debezium.mongodb.Trigger` for CDC.
Connection/auth can succeed, and MongoDB change streams work when tested directly with `mongosh`, but Kestra trigger is unstable and often cannot consume events correctly.

**Environment**
- Kestra: (please fill exact version)
- Plugin: `io.kestra.plugin.debezium.mongodb` (no newer plugin release available currently)
- Debezium in logs: `3.3.1.Final`
- MongoDB: Replica Set `rs0` on K8s

**Trigger config (example)**
```yaml
triggers:
- id: customenote_deidentify_cdc
type: io.kestra.plugin.debezium.mongodb.Trigger
snapshotMode: NO_DATA
connectionString: "{{ kv('MONGODB_URI') }}"
includedDatabases: test
includedCollections: test.encounters
deleted: DROP
stateName: "test_customnote_cdc"
maxWait: PT30S
maxDuration: PT5M
```

**Observed problems**

1. Trigger often completes with no events:
- `Debezium ended successfully ... completed normally`
- `Found '0' messages`
- even when real updates are happening

2. Frequent runtime failure in worker:
- `Engine has been already shut down`
- during trigger evaluation/start-stop cycle

3. Critical processing failure after receiving events:
- `IllegalStateException: DebeziumOpenLineageEmitter not initialized ... Call init() first`
- connector stops with `ConnectException: Error while processing event at offset ...`

4. Another recurring error:
- `Unable to register the MBean 'debezium.mongodb:type=connector-metrics,...'`
- caused by `InterruptedException: sleep interrupted`

- kestra container error logs:
```
04:32:34.417 ERROR tra_-change-event-source-coordinator io.debezium.pipeline.ErrorHandler Producer failure
java.lang.RuntimeException: Unable to register the MBean 'debezium.mongodb:type=connector-metrics,context=snapshot,server=kestra_,task=0'
at io.debezium.pipeline.JmxUtils.registerMXBean(JmxUtils.java:68)
at io.debezium.metrics.Metrics.register(Metrics.java:65)
at io.debezium.pipeline.ChangeEventSourceCoordinator.lambda$start$0(ChangeEventSourceCoordinator.java:135)
at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Unknown Source)
at java.base/java.util.concurrent.FutureTask.run(Unknown Source)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source)
at java.base/java.lang.Thread.run(Unknown Source)
Caused by: java.lang.InterruptedException: sleep interrupted
at java.base/java.lang.Thread.sleepNanos0(Native Method)
at java.base/java.lang.Thread.sleepNanos(Unknown Source)
at java.base/java.lang.Thread.sleep(Unknown Source)
at io.debezium.util.Metronome$1.pause(Metronome.java:57)
at io.debezium.pipeline.JmxUtils.registerMXBean(JmxUtils.java:59)
... 7 common frames omitted
03:23:53.007 ERROR tra_-change-event-source-coordinator io.debezium.pipeline.ErrorHandler Producer failure
org.apache.kafka.connect.errors.ConnectException: Error while processing event at offset {sec=1773717833, ord=2, resume_token=gQAAAAJfZGF0YQBxAAAAODI2OUI4Qzk0OTAwMDAwMDAyMkIwMjJDMDEwMDI5NkU1QTEwMDQ1OTNEQURENDI2RjQ0M0YzODgxQkE4RTk0RjZDRDUxRTQ2NjQ1RjY5NjQwMDY0NjlCOEM5NDk1MDQxNTcwMDA4MTcyMjVEMDAwNAAA}
at io.debezium.pipeline.EventDispatcher.handleEventProcessingFailure(EventDispatcher.java:355)
at io.debezium.pipeline.EventDispatcher.dispatchDataChangeEvent(EventDispatcher.java:347)
at io.debezium.connector.mongodb.MongoDbStreamingChangeEventSource.dispatchChangeEvent(MongoDbStreamingChangeEventSource.java:163)
at io.debezium.connector.mongodb.MongoDbStreamingChangeEventSource.lambda$processChangeStreamDocument$5(MongoDbStreamingChangeEventSource.java:148)
at io.debezium.connector.mongodb.MongoDbStreamingChangeEventSource.errorHandled(MongoDbStreamingChangeEventSource.java:178)
at io.debezium.connector.mongodb.MongoDbStreamingChangeEventSource.lambda$processChangeStreamDocument$6(MongoDbStreamingChangeEventSource.java:148)
at java.base/java.util.Optional.map(Unknown Source)
at io.debezium.connector.mongodb.MongoDbStreamingChangeEventSource.processChangeStreamDocument(MongoDbStreamingChangeEventSource.java:148)
at io.debezium.connector.mongodb.MongoDbStreamingChangeEventSource.lambda$readChangeStream$1(MongoDbStreamingChangeEventSource.java:113)
at java.base/java.util.Optional.map(Unknown Source)
at io.debezium.connector.mongodb.MongoDbStreamingChangeEventSource.readChangeStream(MongoDbStreamingChangeEventSource.java:113)
at io.debezium.connector.mongodb.MongoDbStreamingChangeEventSource.lambda$execute$0(MongoDbStreamingChangeEventSource.java:85)
at io.debezium.connector.mongodb.connection.MongoDbConnection.lambda$execute$0(MongoDbConnection.java:89)
at io.debezium.connector.mongodb.connection.MongoDbConnection.execute(MongoDbConnection.java:105)
at io.debezium.connector.mongodb.connection.MongoDbConnection.execute(MongoDbConnection.java:88)
at io.debezium.connector.mongodb.MongoDbStreamingChangeEventSource.execute(MongoDbStreamingChangeEventSource.java:84)
at io.debezium.connector.mongodb.MongoDbStreamingChangeEventSource.execute(MongoDbStreamingChangeEventSource.java:37)
at io.debezium.pipeline.ChangeEventSourceCoordinator.streamEvents(ChangeEventSourceCoordinator.java:329)
at io.debezium.pipeline.ChangeEventSourceCoordinator.executeChangeEventSources(ChangeEventSourceCoordinator.java:207)
at io.debezium.pipeline.ChangeEventSourceCoordinator.lambda$start$0(ChangeEventSourceCoordinator.java:147)
at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Unknown Source)
at java.base/java.util.concurrent.FutureTask.run(Unknown Source)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source)
at java.base/java.lang.Thread.run(Unknown Source)
Caused by: java.lang.IllegalStateException: DebeziumOpenLineageEmitter not initialized for connector ConnectorContext[connectorLogicalName=kestra_, connectorName=mongodb, taskId=0, version=3.3.1.Final, config={connector.class=io.debezium.connector.mongodb.MongoDbConnector, collection.include.list=empathia.attachments, mongodb.connection.string=mongodb://cdc_user:2ZmFuYS1tb25n1@mongodb-0.mongodb-headless.infra.svc.cluster.local:27017/admin?replicaSet=rs0, record.processing.shutdown.timeout.ms=1000, capture.mode=change_streams_update_full_with_pre_image, record.processing.order=ORDERED, tombstones.on.delete=false, topic.prefix=kestra_, offset.storage.file.filename=/tmp/6Ef1k9JbyhbAToTq8h4cbK/offsets.dat, record.processing.threads=, errors.retry.delay.initial.ms=300, value.converter=org.apache.kafka.connect.json.JsonConverter, key.converter=org.apache.kafka.connect.json.JsonConverter, offset.storage=org.apache.kafka.connect.storage.FileOffsetBackingStore, database.server.name=kestra, offset.flush.timeout.ms=5000, errors.retry.delay.max.ms=10000, offset.flush.interval.ms=1000, key.converter.schemas.enable=false, internal.task.management.timeout.ms=40000, record.processing.with.serial.consumer=false, errors.max.retries=-1, name=engine, value.converter.schemas.enable=false, database.include.list=empathia, snapshot.mode=no_data}]. Call init() first.
at io.debezium.openlineage.DebeziumOpenLineageEmitter.getEmitter(DebeziumOpenLineageEmitter.java:176)
at io.debezium.openlineage.DebeziumOpenLineageEmitter.emit(DebeziumOpenLineageEmitter.java:153)
at io.debezium.connector.mongodb.MongoDbSchema.lambda$schemaFor$0(MongoDbSchema.java:97)
at java.base/java.util.concurrent.ConcurrentHashMap.computeIfAbsent(Unknown Source)
at io.debezium.connector.mongodb.MongoDbSchema.schemaFor(MongoDbSchema.java:74)
at io.debezium.connector.mongodb.MongoDbSchema.schemaFor(MongoDbSchema.java:38)
at io.debezium.pipeline.EventDispatcher.dispatchDataChangeEvent(EventDispatcher.java:287)
... 23 common frames omitted
```

**Important validation**
- Using the same Mongo user/connection in `mongosh`, `db.collection.watch()` receives update events correctly.
- So MongoDB CDC itself is working; issue appears in Kestra trigger runtime lifecycle / Debezium integration path.

**Expected behavior**
- Trigger should consume MongoDB change events reliably in `NO_DATA` mode without crashing.
- No OpenLineage emitter initialization failure when processing first events.
- No MBean registration interruption due to engine lifecycle race.

**Actual behavior**
- Either no events (`0 messages`) or connector crashes with OpenLineage/MBean-related exceptions.

**Request**
- Please confirm if this is a known issue with current Kestra Debezium integration (Debezium 3.3.1 path).
- Provide recommended workaround/config to fully disable problematic OpenLineage/JMX paths for this trigger.
- Provide plugin release/patch plan for a stable MongoDB Debezium Trigger runtime.

### Environment

- Kestra Version: v1.3.0

Contributor guide

No contributing guide indexed for this repository

Research direction

Start with the io.kestra.plugin.debezium.mongodb.Trigger configuration and its runtime lifecycle, reproducing the provided MongoDB Replica Set setup with snapshotMode NO_DATA. Compare trigger behavior with mongosh db.collection.watch(); done means change events are consumed reliably without the reported 0-message, OpenLineage initialization, engine shutdown, or MBean interruption errors.

Written by the indexing model from the issue text.

Assessment

Tech stack
java, mongodb
Domain
backend, databases
Issue type
Bug
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.