[ 
https://issues.apache.org/jira/browse/FLINK-8093?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18119728#comment-18119728
 ] 

Seungki Kim commented on FLINK-8093:
------------------------------------

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)

Reply via email to