A) With the Consumer (confluent-kafka-go and librdkafka from HEAD), there is a Segfault, when partitions are added to a topic.
B) Besides Segfaulting, the consumers do not pick up the new topics, on there own (without changing the consumer group).
Is B) expected behavior? Can we fix A)?
Thank's for your work on the Kafka libraries.
On the broker, create a topic with 2 partitions:
$KAFKA_HOME/bin/kafka-topics.sh --zookeeper localhost:2181 --topic topic --create --replication-factor 1 --partitions 2
Run 2 Consumers with:
confluent-kafka-go/examples/consumer_channel_example/consumer_channel_example.go
On the broker, change the topic configuration with:
$KAFKA_HOME/bin/kafka-topics.sh --zookeeper localhost:2181 --topic topic --alter --partitions 4
Observe 1)
No Consumer picked up the new partitions
Stop 1 Consumer with
Observe 2)
The other Consumer crashes with a Segfault.
Before the Segfault it logs "PARTCNT|[thrd:main]: Topic topic partition count changed from 2 to 4"
Please provide the following information:
Thank you for the report, will investigate!
@edenhill Are librdkafka and confluent-kafka-go supposed to rebalance the partitions using the new sticky strategy introduced in the 0.11?
I am referring to the behavior introduced in the Java Consumer Client (https://github.com/apache/kafka/blob/trunk/clients/src/main/java/org/apache/kafka/clients/consumer/StickyAssignor.java)
Maybe we're running into the same issue. The SEGFAULT happens on the same librdkafka path I believe however I can't reproduce it with the steps described above and we're not performing any topic alter operations, we just create a topic, produce and consume it and then delete it. We use the function-based consumer.
We use the following versions:
Logs:
%7|1499689330.544|SEND|rdkafka#consumer-288| [thrd:kc2.docker:9092/bootstrap]: kc2.docker:9092/2: Sent GroupCoordinatorRequest (v0, 29 bytes @ 0, CorrId 2) [7/3689]
%7|1499689330.544|STATE|rdkafka#consumer-288| [thrd:kc2.docker:9092/bootstrap]: kc2.docker:9092/2: Broker changed state UPDATE -> UP
%7|1499689330.544|BROADCAST|rdkafka#consumer-288| [thrd:kc2.docker:9092/bootstrap]: Broadcasting state change
%7|1499689330.545|RECV|rdkafka#consumer-288| [thrd:kc2.docker:9092/bootstrap]: kc2.docker:9092/2: Received GroupCoordinatorResponse (v0, 22 bytes, CorrId 2, rtt 1.50ms)
%7|1499689330.546|CGRPCOORD|rdkafka#consumer-288| [thrd:main]: kc2.docker:9092/2: Group "1678b6" coordinator is kc1.docker:9092 id 1
%7|1499689330.546|CGRPCOORD|rdkafka#consumer-288| [thrd:main]: Group "1678b6" changing coordinator -1 -> 1
%7|1499689330.546|CGRPSTATE|rdkafka#consumer-288| [thrd:main]: Group "1678b6" changed state wait-broker -> wait-broker-transport (v1, join-state init)
%7|1499689330.546|BROADCAST|rdkafka#consumer-288| [thrd:main]: Broadcasting state change
%7|1499689330.546|CGRPSTATE|rdkafka#consumer-288| [thrd:main]: Group "1678b6" changed state wait-broker-transport -> up (v1, join-state init)
%7|1499689330.546|BROADCAST|rdkafka#consumer-288| [thrd:main]: Broadcasting state change
%7|1499689330.546|JOIN|rdkafka#consumer-288| [thrd:main]: Group "1678b6": join with 0 (1) subscribed topic(s)
%7|1499689330.546|METADATA|rdkafka#consumer-288| [thrd:main]: Hinted cache of 1/1 topic(s) being queried
%7|1499689330.546|CGRPMETADATA|rdkafka#consumer-288| [thrd:main]: consumer join: metadata for subscription only available for 0/1 topics (-1ms old)
%7|1499689330.546|METADATA|rdkafka#consumer-288| [thrd:main]: kc2.docker:9092/2: Request metadata for 1 topic(s): consumer join
%7|1499689330.546|JOIN|rdkafka#consumer-288| [thrd:main]: Group "1678b6": postponing join until up-to-date metadata is available
%7|1499689330.546|SEND|rdkafka#consumer-288| [thrd:kc2.docker:9092/bootstrap]: kc2.docker:9092/2: Sent MetadataRequest (v2, 35 bytes @ 0, CorrId 3)
%7|1499689330.546|RECV|rdkafka#consumer-288| [thrd:kc2.docker:9092/bootstrap]: kc2.docker:9092/2: Received MetadataResponse (v2, 233 bytes, CorrId 3, rtt 0.68ms)
%7|1499689330.546|METADATA|rdkafka#consumer-288| [thrd:main]: kc2.docker:9092/2: ===== Received metadata (for 1 requested topics): consumer join =====
%7|1499689330.546|METADATA|rdkafka#consumer-288| [thrd:main]: kc2.docker:9092/2: ClusterId: PofQnMJSSGqcKlC77DCxmQ, ControllerId: 1
%7|1499689330.546|METADATA|rdkafka#consumer-288| [thrd:main]: kc2.docker:9092/2: 2 brokers, 1 topics
%7|1499689330.546|METADATA|rdkafka#consumer-288| [thrd:main]: kc2.docker:9092/2: Broker #0/2: kc2.docker:9092 NodeId 2
%7|1499689330.547|METADATA|rdkafka#consumer-288| [thrd:main]: kc2.docker:9092/2: Broker #1/2: kc1.docker:9092 NodeId 1
%7|1499689330.547|METADATA|rdkafka#consumer-288| [thrd:main]: kc2.docker:9092/2: Topic #0/1: r-9a4e4b with 4 partitions
%7|1499689330.547|METADATA|rdkafka#consumer-288| [thrd:main]: kc2.docker:9092/2: 1/1 requested topic(s) seen in metadata
%7|1499689330.547|SUBSCRIPTION|rdkafka#consumer-288| [thrd:main]: Group "1678b6": effective subscription list changed from 0 to 1 topic(s):
%7|1499689330.547|SUBSCRIPTION|rdkafka#consumer-288| [thrd:main]: Topic r-9a4e4b with 4 partition(s)
%7|1499689330.547|JOIN|rdkafka#consumer-288| [thrd:main]: Group "1678b6": join with 1 (1) subscribed topic(s)
%7|1499689330.547|CGRPMETADATA|rdkafka#consumer-288| [thrd:main]: consumer join: metadata for subscription is up to date (0ms old)
%7|1499689330.547|CGRPJOINSTATE|rdkafka#consumer-288| [thrd:main]: Group "1678b6" changed join state init -> wait-join (v1, state up)
%7|1499689330.547|SEND|rdkafka#consumer-288| [thrd:kc1.docker:9092/bootstrap]: kc1.docker:9092/1: Sent JoinGroupRequest (v0, 116 bytes @ 0, CorrId 3)
%7|1499689330.547|BROADCAST|rdkafka#consumer-288| [thrd:kc1.docker:9092/bootstrap]: Broadcasting state change
%7|1499689330.548|RECV|rdkafka#consumer-288| [thrd:kc1.docker:9092/bootstrap]: kc1.docker:9092/1: Received JoinGroupResponse (v0, 179 bytes, CorrId 3, rtt 1.29ms)
%7|1499689330.548|JOINGROUP|rdkafka#consumer-288| [thrd:main]: JoinGroup response: GenerationId 1, Protocol range, LeaderId rdkafka-42fba3ad-4ff6-44b3-a686-0993a6f28703 (me), my MemberId rdkafka-42fba3ad-4ff6-44b3-a686-0993a6f28703, 1 members in group: (no error)
%7|1499689330.548|MEMBERID|rdkafka#consumer-288| [thrd:main]: Group "1678b6": updating member id "" -> "rdkafka-42fba3ad-4ff6-44b3-a686-0993a6f28703"
%7|1499689330.548|JOINGROUP|rdkafka#consumer-288| [thrd:main]: Elected leader for group "1678b6" with 1 member(s)
%7|1499689330.548|CGRPJOINSTATE|rdkafka#consumer-288| [thrd:main]: Group "1678b6" changed join state wait-join -> wait-metadata (v1, state up)
%7|1499689330.548|METADATA|rdkafka#consumer-288| [thrd:main]: kc1.docker:9092/1: Request metadata for 1 topic(s): partition assignor
%7|1499689330.549|SEND|rdkafka#consumer-288| [thrd:kc1.docker:9092/bootstrap]: kc1.docker:9092/1: Sent MetadataRequest (v2, 35 bytes @ 0, CorrId 4)
%7|1499689330.550|RECV|rdkafka#consumer-288| [thrd:kc1.docker:9092/bootstrap]: kc1.docker:9092/1: Received MetadataResponse (v2, 97 bytes, CorrId 4, rtt 1.01ms)
%7|1499689330.550|METADATA|rdkafka#consumer-288| [thrd:main]: kc1.docker:9092/1: ===== Received metadata (for 1 requested topics): partition assignor =====
%7|1499689330.550|METADATA|rdkafka#consumer-288| [thrd:main]: kc1.docker:9092/1: ClusterId: PofQnMJSSGqcKlC77DCxmQ, ControllerId: 1
%7|1499689330.550|METADATA|rdkafka#consumer-288| [thrd:main]: kc1.docker:9092/1: 2 brokers, 1 topics
%7|1499689330.550|METADATA|rdkafka#consumer-288| [thrd:main]: kc1.docker:9092/1: Broker #0/2: kc2.docker:9092 NodeId 2
%7|1499689330.550|METADATA|rdkafka#consumer-288| [thrd:main]: kc1.docker:9092/1: Broker #1/2: kc1.docker:9092 NodeId 1
%7|1499689330.550|METADATA|rdkafka#consumer-288| [thrd:main]: kc1.docker:9092/1: Topic #0/1: r-9a4e4b with 0 partitions: Broker: Unknown topic or partition
%7|1499689330.550|METADATA|rdkafka#consumer-288| [thrd:main]: kc1.docker:9092/1: 1/1 requested topic(s) seen in metadata
%7|1499689330.550|SUBSCRIPTION|rdkafka#consumer-288| [thrd:main]: Group "1678b6": no topics in metadata matched subscription
%7|1499689330.550|SUBSCRIPTION|rdkafka#consumer-288| [thrd:main]: Group "1678b6": effective subscription list changed from 1 to 0 topic(s):
make: *** [spawn] Error 139
Here's a gdb session in case it helps:
(gdb) bt
#0 strlen () at ../sysdeps/x86_64/strlen.S:106
#1 0x00007fe087324e0f in rd_kafkap_str_cmp_str2 (str=0x0, b=0x7fe078002520) at rdkafka_proto.h:302
#2 0x00007fe0873264db in rd_kafka_assignor_cmp_str (_a=0x0, _b=0x7fe0780024d0) at rdkafka_assignor.c:388
#3 0x00007fe08732bc11 in rd_list_find (rl=0x7fe078003898, match=0x0, cmp=0x7fe0873264a4 <rd_kafka_assignor_cmp_str>) at rdlist.c:222
#4 0x00007fe08732650e in rd_kafka_assignor_find (rk=0x7fe0780034a0, protocol=0x0) at rdkafka_assignor.c:399
#5 0x00007fe087325e3f in rd_kafka_assignor_run (rkcg=0x7fe078002570, protocol_name=0x0, metadata=0x7fe058001560, members=0x0, member_cnt=0, errstr=0x7fe069ff7420 "", errstr_size=512) at rdkafka_assignor.c:287
#6 0x00007fe08730df2c in rd_kafka_cgrp_assignor_run (rkcg=0x7fe078002570, protocol_name=0x0, err=RD_KAFKA_RESP_ERR_NO_ERROR, metadata=0x7fe058001560, members=0x0, member_cnt=0) at rdkafka_cgrp.c:663
#7 0x00007fe08730e1c8 in rd_kafka_cgrp_assignor_handle_Metadata_op (rk=0x7fe0780034a0, rkq=0x7fe069ff7770, rko=0x7fe058000ee0) at rdkafka_cgrp.c:717
#8 0x00007fe0872fb2c1 in rd_kafka_op_call (rk=0x7fe0780034a0, rkq=0x7fe069ff7770, rko=0x7fe058000ee0) at rdkafka_op.c:492
#9 0x00007fe0872fb546 in rd_kafka_op_handle_std (rk=0x7fe0780034a0, rkq=0x7fe069ff7770, rko=0x7fe058000ee0, cb_type=1) at rdkafka_op.c:586
#10 0x00007fe0872fb5e1 in rd_kafka_op_handle (rk=0x7fe0780034a0, rkq=0x7fe069ff7770, rko=0x7fe058000ee0, cb_type=RD_KAFKA_Q_CB_CALLBACK, opaque=0x0, callback=0x0) at rdkafka_op.c:621
#11 0x00007fe0872f8600 in rd_kafka_q_serve (rkq=0x7fe078002380, timeout_ms=981, max_cnt=0, cb_type=RD_KAFKA_Q_CB_CALLBACK, callback=0x0, opaque=0x0) at rdkafka_queue.c:467
#12 0x00007fe0872c6685 in rd_kafka_thread_main (arg=0x7fe0780034a0) at rdkafka.c:1227
#13 0x00007fe08732c354 in _thrd_wrapper_function (aArg=0x7fe078002ad0) at tinycthread.c:624
#14 0x00007fe08709a494 in start_thread (arg=0x7fe069ffb700) at pthread_create.c:333
#15 0x00007fe086ddcaff in clone () at ../sysdeps/unix/sysv/linux/x86_64/clone.S:97
Let me know if you want me to upload the core file somewhere.
This will be fixed in the next release of librdkafka (0.11.0) which is imminent:
https://github.com/edenhill/librdkafka/issues/1193
@edenhill Unfortunately, this persists with librdkafka 0.11 which contains the commit supposed to fix this (https://github.com/edenhill/librdkafka/commit/c59e479caf28ab481178e14c092b8ab5e5b77ece).
I think the backtrace is identical to my previous one, which was with 0.9.5, RC1 or RC2 (can't remember):
(gdb) bt
#0 strlen () at ../sysdeps/x86_64/strlen.S:106
#1 0x00007fc6ab177fe6 in rd_kafkap_str_cmp_str2 (str=0x0, b=0x133d330) at rdkafka_proto.h:302
#2 0x00007fc6ab1796b2 in rd_kafka_assignor_cmp_str (_a=0x0, _b=0x133cd40) at rdkafka_assignor.c:388
#3 0x00007fc6ab17edeb in rd_list_find (rl=0x133c4d8, match=0x0, cmp=0x7fc6ab17967b <rd_kafka_assignor_cmp_str>) at rdlist.c:225
#4 0x00007fc6ab1796e5 in rd_kafka_assignor_find (rk=0x133c0e0, protocol=0x0) at rdkafka_assignor.c:399
#5 0x00007fc6ab179016 in rd_kafka_assignor_run (rkcg=0x133cea0, protocol_name=0x0, metadata=0x7fc694000d90, members=0x0, member_cnt=0, errstr=0x7fc68fde4420 "", errstr_size=512) at rdkafka_assignor.c:287
#6 0x00007fc6ab161136 in rd_kafka_cgrp_assignor_run (rkcg=0x133cea0, protocol_name=0x0, err=RD_KAFKA_RESP_ERR_NO_ERROR, metadata=0x7fc694000d90, members=0x0, member_cnt=0) at rdkafka_cgrp.c:663
#7 0x00007fc6ab1613d2 in rd_kafka_cgrp_assignor_handle_Metadata_op (rk=0x133c0e0, rkq=0x7fc68fde4770, rko=0x7fc6940012d0) at rdkafka_cgrp.c:717
#8 0x00007fc6ab14e4cb in rd_kafka_op_call (rk=0x133c0e0, rkq=0x7fc68fde4770, rko=0x7fc6940012d0) at rdkafka_op.c:492
#9 0x00007fc6ab14e750 in rd_kafka_op_handle_std (rk=0x133c0e0, rkq=0x7fc68fde4770, rko=0x7fc6940012d0, cb_type=1) at rdkafka_op.c:586
#10 0x00007fc6ab14e7eb in rd_kafka_op_handle (rk=0x133c0e0, rkq=0x7fc68fde4770, rko=0x7fc6940012d0, cb_type=RD_KAFKA_Q_CB_CALLBACK, opaque=0x0, callback=0x0) at rdkafka_op.c:621
#11 0x00007fc6ab14b80a in rd_kafka_q_serve (rkq=0x133d250, timeout_ms=988, max_cnt=0, cb_type=RD_KAFKA_Q_CB_CALLBACK, callback=0x0, opaque=0x0) at rdkafka_queue.c:467
#12 0x00007fc6ab1196b5 in rd_kafka_thread_main (arg=0x133c0e0) at rdkafka.c:1227
#13 0x00007fc6ab17f52e in _thrd_wrapper_function (aArg=0x1339d60) at tinycthread.c:624
#14 0x00007fc6aaeed494 in start_thread (arg=0x7fc68fde8700) at pthread_create.c:333
#15 0x00007fc6aac2faff in clone () at ../sysdeps/unix/sysv/linux/x86_64/clone.S:97
So I think the path is a different one, we don't have a JoinGroup response but a Metadata request that receives an empty protocol (ie. rd_kafka_cgrp_assignor_handle_Metadata_op() calls rd_kafka_cgrp_assignor_run()).
Ouch, ok, will investigate.
I have also run into this issue in v0.11 using the original repro steps described. After some debugging I found that if a metadata update is received while the Kafka consumer group join state is set to RD_KAFKA_CGRP_JOIN_STATE_WAIT_METADATA, the metadata parser determines that the consumer needs to rejoin the group causing rkcg->rkcg_group_leader.protocol to be set to NULL in the following code path:
#0 rd_kafka_cgrp_group_leader_reset (rkcg=0x6fe130) at rdkafka_cgrp.c:2388
#1 0x0000000000459d66 in rd_kafka_cgrp_rejoin (rkcg=0x6fe130) at rdkafka_cgrp.c:1185
#2 0x000000000045ecd2 in rd_kafka_cgrp_metadata_update_check (rkcg=0x6fe130, do_join=1) at rdkafka_cgrp.c:3083
#3 0x000000000047e0bc in rd_kafka_parse_Metadata (rkb=0x701520, request=0x7038b0, rkbuf=0x700860) at rdkafka_metadata.c:551
#4 0x000000000044bb18 in rd_kafka_handle_Metadata (rk=0x6fd5a0, rkb=0x701520, err=RD_KAFKA_RESP_ERR_NO_ERROR, rkbuf=0x700860,
request=0x7038b0, opaque=0x7024b0) at rdkafka_request.c:1293
#5 0x000000000043bac3 in rd_kafka_buf_callback (rk=0x6fd5a0, rkb=0x701520, err=RD_KAFKA_RESP_ERR_NO_ERROR, response=0x700860,
request=0x7038b0) at rdkafka_buf.c:421
#6 0x000000000043b897 in rd_kafka_buf_handle_op (rko=0x703bf0, err=RD_KAFKA_RESP_ERR_NO_ERROR) at rdkafka_buf.c:362
#7 0x00000000004410d6 in rd_kafka_op_handle_std (rk=0x6fd5a0, rkq=0x7ffff7426930, rko=0x703bf0, cb_type=1) at rdkafka_op.c:588
#8 0x0000000000441152 in rd_kafka_op_handle (rk=0x6fd5a0, rkq=0x7ffff7426930, rko=0x703bf0, cb_type=RD_KAFKA_Q_CB_CALLBACK,
opaque=0x0, callback=0x0) at rdkafka_op.c:621
#9 0x000000000043e01c in rd_kafka_q_serve (rkq=0x6fde80, timeout_ms=0, max_cnt=0, cb_type=RD_KAFKA_Q_CB_CALLBACK, callback=0x0,
opaque=0x0) at rdkafka_queue.c:467
#10 0x000000000040b111 in rd_kafka_thread_main (arg=0x6fd5a0) at rdkafka.c:1227
#11 0x00000000004745c1 in _thrd_wrapper_function (aArg=0x6fe6b0) at tinycthread.c:624
#12 0x00007ffff7bc69ca in start_thread () from /lib/libpthread.so.0
#13 0x00007ffff751645d in clone () from /lib/libc.so.6
rd_kafka_cgrp_rejoin would then invoke rd_kafka_cgrp_join, but since rkcg->rkcg_join_state != RD_KAFKA_CGRP_JOIN_STATE_INIT, it would return without doing anything, leaving the protocol field NULL. This is the last time the protocol field is modified before the segmentation fault occurs.
I've been testing out a patch for edenhill/librdkafka (attached below) where I've moved the call to reset the group leader struct to happen only when new assignments are present or if the consumer is in a state to join a group as defined in rd_kafka_cgrp_join. This seems to fix the segmentation fault issue and I have not yet encountered any regressions. However, I'm new to Kafka's protocol and the librdkafka codebase, so I do not know if this is a correct fix. I hope this provided some more useful information though.
EDIT: Found an alternative and potentially more correct fix. Since there have been detected changes in the metadata that require a re-joining of the group, allowing to make a join request when the join_state is RD_KAFKA_CGRP_JOIN_STATE_WAIT_METADATA will also fix the segmentation fault issue.
diff --git a/src/rdkafka_cgrp.c b/src/rdkafka_cgrp.c
index 3e6185d..1a66afe 100644
--- a/src/rdkafka_cgrp.c
+++ b/src/rdkafka_cgrp.c
@@ -1118,7 +1118,8 @@ static void rd_kafka_cgrp_join (rd_kafka_cgrp_t *rkcg) {
int metadata_age;
if (rkcg->rkcg_state != RD_KAFKA_CGRP_STATE_UP ||
- rkcg->rkcg_join_state != RD_KAFKA_CGRP_JOIN_STATE_INIT)
+ (rkcg->rkcg_join_state != RD_KAFKA_CGRP_JOIN_STATE_INIT &&
+ rkcg->rkcg_join_state != RD_KAFKA_CGRP_JOIN_STATE_WAIT_METADATA))
return;
rd_kafka_dbg(rkcg->rkcg_rk, CGRP, "JOIN",
@matavine Great investigation, absolutely right. :100:
This has now been fixed on librdkafka master.
Most helpful comment
This will be fixed in the next release of librdkafka (0.11.0) which is imminent:
https://github.com/edenhill/librdkafka/issues/1193