[
https://issues.apache.org/jira/browse/FLINK-8093?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18119728#comment-18119728
]
Seungki Kim edited comment on FLINK-8093 at 9/28/26 12:49 AM:
--------------------------------------------------------------
Hi [~nateab], I noticed that PR #226 was auto-closed as stale in June. Are you
still planning to work on this? If not, I'd be happy to take it over and build
on your PR.
Some context from our side: we run several KafkaSinks with parallelism > 1 and
set a fixed `{{{}client.id{}}}` so we can identify our producers on the broker.
Every parallel producer then registers the same id, which produces
`{{{}InstanceAlreadyExistsException{}}}` warnings for the app-info MBean and
makes the per-client broker metrics ambiguous. KafkaSource already solves this
with `{{{}client.id.prefix{}}}`, so we'd like the same on the sink side.
Based on the earlier attempts (#101, #118, #226) and the review feedback on
#101, here's what I have in mind:
# Add `{{{}KafkaSinkOptions.CLIENT_ID_PREFIX{}}}` and
`{{{}KafkaSinkBuilder#setClientIdPrefix{}}}`, mirroring the source. If both the
prefix and {{client.id}} are set, the prefix wins, as it does on the source
side.
# Derive a unique id for every Kafka client the sink creates: the
at-least-once writer producer, the pooled transactional producers, the
committer producer, the AdminClient in `{{{}ExactlyOnceKafkaWriter{}}}`, and
the metadata producer in `{{{}DefaultKafkaSinkContext{}}}`. The format would be
`{{{}<prefix>\-<subtaskId>\-<role>-<counter>{}}}`. Since pooled producers are
reused across transactional ids, the id can't be tied to the transactional id;
the counter guarantees uniqueness among producers that are open at the same
time (the concern raised on #101).
# Testing: unit tests asserting that all generated ids are unique, plus an
ITCase with parallelism > 1 that checks the registered
`{{{}kafka.producer:type=app-info{}}}` MBeans do not collide, instead of
relying only on manual log inspection.
I'd also propose tracking this under this ticket only and closing FLINK-28842
and FLINK-35283 as duplicates. The Table API option
(`{{{}sink.client-id-prefix{}}}`) could be a follow-up.
cc. [~arvid]
was (Author: JIRAUSER314455):
Hi [~nateab], I noticed that PR #226 was auto-closed as stale in June. Are you
still planning to work on this? If not, I'd be happy to take it over and build
on your PR.
Some context from our side: we run several KafkaSinks with parallelism > 1 and
set a fixed `{{{}client.id{}}}` so we can identify our producers on the broker.
Every parallel producer then registers the same id, which produces
`{{{}InstanceAlreadyExistsException{}}}` warnings for the app-info MBean and
makes the per-client broker metrics ambiguous. KafkaSource already solves this
with `{{{}client.id.prefix{}}}`, so we'd like the same on the sink side.
Based on the earlier attempts (#101, #118, #226) and the review feedback on
#101, here's what I have in mind:
# Add `{{{}KafkaSinkOptions.CLIENT_ID_PREFIX{}}}` and
`{{{}KafkaSinkBuilder#setClientIdPrefix{}}}`, mirroring the source. If both the
prefix and {{client.id}} are set, the prefix wins, as it does on the source
side.
# Derive a unique id for every Kafka client the sink creates: the
at-least-once writer producer, the pooled transactional producers, the
committer producer, the AdminClient in `{{{}ExactlyOnceKafkaWriter{}}}`, and
the metadata producer in `{{{}DefaultKafkaSinkContext{}}}`. The format would be
`{{{}<prefix>-<subtaskId>-<role>-<counter>{}}}`. Since pooled producers are
reused across transactional ids, the id can't be tied to the transactional id;
the counter guarantees uniqueness among producers that are open at the same
time (the concern raised on #101).
# Testing: unit tests asserting that all generated ids are unique, plus an
ITCase with parallelism > 1 that checks the registered
`{{{}kafka.producer:type=app-info{}}}` MBeans do not collide, instead of
relying only on manual log inspection.
I'd also propose tracking this under this ticket only and closing FLINK-28842
and FLINK-35283 as duplicates. The Table API option
(`{{{}sink.client-id-prefix{}}}`) could be a follow-up.
cc. [~arvid]
> flink job fail because of kafka producer create fail of
> "javax.management.InstanceAlreadyExistsException"
> ---------------------------------------------------------------------------------------------------------
>
> Key: FLINK-8093
> URL: https://issues.apache.org/jira/browse/FLINK-8093
> Project: Flink
> Issue Type: Bug
> Components: Connectors / Kafka
> Affects Versions: 1.3.2, 1.10.0
> Environment: flink 1.3.2, kafka 0.9.1
> Reporter: dongtingting
> Assignee: Natea Eshetu Beshada
> Priority: Not a Priority
> Labels: auto-deprioritized-critical, auto-deprioritized-major,
> auto-deprioritized-minor, auto-unassigned, pull-request-available, usability
>
> one taskmanager has multiple taskslot, one task fail because of create
> kafkaProducer fail,the reason for create kafkaProducer fail is
> “javax.management.InstanceAlreadyExistsException:
> kafka.producer:type=producer-metrics,client-id=producer-3”。 the detail trace
> is :
> {noformat}
> 2017-11-04 19:41:23,281 INFO org.apache.flink.runtime.taskmanager.Task
> - Source: Custom Source -> Filter -> Map -> Filter -> Sink:
> dp_client_**_log (7/80) (99551f3f892232d7df5eb9060fa9940c) switched from
> RUNNING to FAILED.
> org.apache.kafka.common.KafkaException: Failed to construct kafka producer
> at
> org.apache.kafka.clients.producer.KafkaProducer.<init>(KafkaProducer.java:321)
> at
> org.apache.kafka.clients.producer.KafkaProducer.<init>(KafkaProducer.java:181)
> at
> org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducerBase.getKafkaProducer(FlinkKafkaProducerBase.java:202)
> at
> org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducerBase.open(FlinkKafkaProducerBase.java:212)
> at
> org.apache.flink.api.common.functions.util.FunctionUtils.openFunction(FunctionUtils.java:36)
> at
> org.apache.flink.streaming.api.operators.AbstractUdfStreamOperator.open(AbstractUdfStreamOperator.java:111)
> at
> org.apache.flink.streaming.runtime.tasks.StreamTask.openAllOperators(StreamTask.java:375)
> at
> org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:252)
> at org.apache.flink.runtime.taskmanager.Task.run(Task.java:702)
> at java.lang.Thread.run(Thread.java:745)
> Caused by: org.apache.kafka.common.KafkaException: Error registering mbean
> kafka.producer:type=producer-metrics,client-id=producer-3
> at
> org.apache.kafka.common.metrics.JmxReporter.reregister(JmxReporter.java:159)
> at
> org.apache.kafka.common.metrics.JmxReporter.metricChange(JmxReporter.java:77)
> at
> org.apache.kafka.common.metrics.Metrics.registerMetric(Metrics.java:288)
> at org.apache.kafka.common.metrics.Metrics.addMetric(Metrics.java:255)
> at org.apache.kafka.common.metrics.Metrics.addMetric(Metrics.java:239)
> at
> org.apache.kafka.clients.producer.internals.RecordAccumulator.registerMetrics(RecordAccumulator.java:137)
> at
> org.apache.kafka.clients.producer.internals.RecordAccumulator.<init>(RecordAccumulator.java:111)
> at
> org.apache.kafka.clients.producer.KafkaProducer.<init>(KafkaProducer.java:261)
> ... 9 more
> Caused by: javax.management.InstanceAlreadyExistsException:
> kafka.producer:type=producer-metrics,client-id=producer-3
> at com.sun.jmx.mbeanserver.Repository.addMBean(Repository.java:437)
> at
> com.sun.jmx.interceptor.DefaultMBeanServerInterceptor.registerWithRepository(DefaultMBeanServerInterceptor.java:1898)
> at
> com.sun.jmx.interceptor.DefaultMBeanServerInterceptor.registerDynamicMBean(DefaultMBeanServerInterceptor.java:966)
> at
> com.sun.jmx.interceptor.DefaultMBeanServerInterceptor.registerObject(DefaultMBeanServerInterceptor.java:900)
> at
> com.sun.jmx.interceptor.DefaultMBeanServerInterceptor.registerMBean(DefaultMBeanServerInterceptor.java:324)
> at
> com.sun.jmx.mbeanserver.JmxMBeanServer.registerMBean(JmxMBeanServer.java:522)
> at
> org.apache.kafka.common.metrics.JmxReporter.reregister(JmxReporter.java:157)
> ... 16 more
> {noformat}
> I doubt that task in different taskslot of one taskmanager use different
> classloader, and taskid may be the same in one process。 So this lead to
> create kafkaProducer fail in one taskManager。
> Does anybody encountered the same problem?
--
This message was sent by Atlassian Jira
(v8.20.10#820010)