projects
/
collectd.git
/ blobdiff
commit
grep
author
committer
pickaxe
?
search:
re
summary
|
shortlog
|
log
|
commit
|
commitdiff
|
tree
raw
|
inline
| side by side
Merge pull request #902 from mfournier/write_http-node-blocks
[collectd.git]
/
src
/
write_kafka.c
diff --git
a/src/write_kafka.c
b/src/write_kafka.c
index
ba76d71
..
a2947d1
100644
(file)
--- a/
src/write_kafka.c
+++ b/
src/write_kafka.c
@@
-76,8
+76,13
@@
static int32_t kafka_partition(const rd_kafka_topic_t *rkt,
int32_t partition_cnt, void *p, void *m)
{
u_int32_t key = *((u_int32_t *)keydata );
int32_t partition_cnt, void *p, void *m)
{
u_int32_t key = *((u_int32_t *)keydata );
+ u_int32_t target = key % partition_cnt;
+ int32_t i = partition_cnt;
- return key % partition_cnt;
+ while (--i > 0 && !rd_kafka_topic_partition_available(rkt, target)) {
+ target = (target + 1) % partition_cnt;
+ }
+ return target;
}
static int kafka_write(const data_set_t *ds, /* {{{ */
}
static int kafka_write(const data_set_t *ds, /* {{{ */
@@
-240,7
+245,7
@@
static void kafka_config_topic(rd_kafka_conf_t *conf, oconfig_item_t *ci) /* {{{
goto errout;
}
key = child->values[0].value.string;
goto errout;
}
key = child->values[0].value.string;
- val = child->values[
0
].value.string;
+ val = child->values[
1
].value.string;
ret = rd_kafka_topic_conf_set(tctx->conf,key, val,
errbuf, sizeof(errbuf));
if (ret != RD_KAFKA_CONF_OK) {
ret = rd_kafka_topic_conf_set(tctx->conf,key, val,
errbuf, sizeof(errbuf));
if (ret != RD_KAFKA_CONF_OK) {