Aiven-Open / Aiven-Open/opensearch-connector-for-apache-kafka

"behavior.on.null.values": "delete" doesn't work properly

Đang mở
#203 1 bình luận 0 reaction 0 người được giao Xem trên GitHub
Ngôn ngữ chính
Java
Star
95
Fork
61
Merge trung bình
11 giờ 52 phút
Pull request đã merge (30 ngày)
7

Mô tả

I'm working with opensearch-connector-for-apache-kafka-3.0.0.
The "behavior.on.null.values": "delete" doesn't work properly.
Sometime delete with success.
Othertime it doesn't delete and there are several duplicated elements in the indexes.
This is my opensearch sink configuration:

{
"name": "open-sink-yyy-user-profiling-xxx",
"config": {
"connector.class": "io.aiven.kafka.connect.opensearch.OpensearchSinkConnector",
"type.name": "kafkaconnect-yyy-crsm",
"behavior.on.null.values": "delete",
"connection.password": "",
"topics": "yyy-attributes-category-xxx,yyy-business-user-profiles-xxx,yyy-user-roles-xxx,yyy-business-users-xxx,yyy-offering-business-users-xxx,yyy-users-xxx,yyy-offering-users-xxx,yyy-user-attributes-xxx,yyy-user-relationships-xxx,yyy-user-invitations-xxx,yyy-users-devices-xxx,yyy-iam-group-mappings-xxx",
"tasks.max": "10",
"batch.size": "500",
"connection.timeout.ms": "30000",
"connection.username": "",
"max.retries": "12",
"retry.backoff.ms": "1000",
"schema.ignore": "true",
"value.converter.schemas.enable": "false",
"name": "open-sink-yyy-user-profiling-xxx",
"errors.tolerance": "all",
"connection.url": "http://dcpp-opensearch-cluster:9201",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"read.timeout.ms": "30000",
"key.converter": "org.apache.kafka.connect.storage.StringConverter",
"behavior.on.version.conflict": "warn"
},

We tried also without this setting "behavior.on.version.conflict": "warn"

But in that case the connectors fails and doesnt't work for any operation with this kind of exception

"trace": "org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception.
at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:560)
at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:321)
at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:224)
at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:192)
at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:175)
at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:219)
at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
at java.base/java.lang.Thread.run(Thread.java:834)
nCaused by: org.apache.kafka.connect.errors.ConnectException: Failed to bulk processing after total of 1 attempt(s)
at io.aiven.kafka.connect.opensearch.RetryUtil.callWithRetry(RetryUtil.java:137)at io.aiven.kafka.connect.opensearch.BulkProcessor$BulkTask.execute
(BulkProcessor.java:367)
at io.aiven.kafka.connect.opensearch.BulkProcessor$BulkTask.call(BulkProcessor.java:356)
at io.aiven.kafka.connect.opensearch.BulkProcessor$BulkTask.call(BulkProcessor.java:337)...
4 more\nCaused by: org.apache.kafka.connect.errors.ConnectException: One of the item in the bulk response failed. Reason:
[yyy-user-invitations-xxx/0FNsOYA_R2-Qx7JM04k6Fw][[yyy-user-invitations-xxx][0]]
OpenSearchException[OpenSearch exception [type=version_conflict_engine_exception,
reason=[yyy-1f585694-df20-4a97-bb29-0e73fb5a4832]: version conflict, current version [3] is higher or equal to the one provided [0]]]
at io.aiven.kafka.connect.opensearch.BulkProcessor$BulkTask.handleVersionConflict
(BulkProcessor.java:429)at io.aiven.kafka.connect.opensearch.BulkProcessor$BulkTask.lambda$execute$0(BulkProcessor.java:382)
at io.aiven.kafka.connect.opensearch.RetryUtil.callWithRetry(RetryUtil.java:119)... 7 more\n"

Please can you help me?
Regards
Gabriella

Hướng dẫn đóng góp

Mở hướng dẫn đóng góp

Đánh giá

Issue này chưa được đánh giá.

Nhận issue mới trong hộp thư của bạn

Bản tóm tắt ngắn những issue GitHub phù hợp với người mới.