Madhuvishy has uploaded a new change for review.
https://gerrit.wikimedia.org/r/232408
Change subject: [WIP] Change kafka writer to use pykafka Producer
......................................................................
[WIP] Change kafka writer to use pykafka Producer
Since we moved the kafka reader to use pykafka's BalancedConsumer
instead of python-kafka, also moving the writer code to use
the same library for consistency
Bug: T109244
Change-Id: I5a1b801316786ba017836946ed6e96258187695b
---
M server/eventlogging/handlers.py
1 file changed, 12 insertions(+), 60 deletions(-)
git pull ssh://gerrit.wikimedia.org:29418/mediawiki/extensions/EventLogging
refs/changes/08/232408/1
diff --git a/server/eventlogging/handlers.py b/server/eventlogging/handlers.py
index 39fd7ae..51641e7 100644
--- a/server/eventlogging/handlers.py
+++ b/server/eventlogging/handlers.py
@@ -18,13 +18,7 @@
import json
from functools import partial
-from kafka import KafkaClient
-from kafka import KeyedProducer
-from kafka import SimpleProducer
-from kafka.producer.base import Producer
-from kafka.common import KafkaTimeoutError
-from pykafka import KafkaClient as PyKafkaClient
-from pykafka import BalancedConsumer
+from pykafka import KafkaClient, BalancedConsumer, Producer
import logging
import logging.handlers
@@ -84,7 +78,6 @@
@writes('kafka')
def kafka_writer(
path,
- producer='simple',
topic='eventlogging_%(schema)s',
key='%(schema)s_%(revision)s',
blacklist=None,
@@ -106,8 +99,6 @@
path - URI path should be comma separated Kafka Brokers.
e.g. kafka01:9092,kafka02:9092,kafka03:9092
- producer - Either 'keyed' or 'simple'. Default: 'simple'.
-
topic - Python format string topic name.
If the incoming event is a dict (not a raw string)
topic will be interpolated against event. I.e.
@@ -128,7 +119,9 @@
"""
# Brokers should be in the uri path
- brokers = path.strip('/')
+ # path.strip returns type 'unicode' and pykafka expects a string,
+ # so converting unicode to str
+ brokers = path.strip('/').encode('ascii', 'ignore')
# remove non Kafka Producer args from kafka_consumer_args
kafka_producer_args = {
@@ -136,24 +129,11 @@
if k in inspect.getargspec(Producer.__init__).args
}
- # Use async producer by default
- if 'async' not in kafka_producer_args:
- kafka_producer_args['async'] = True
-
- kafka = KafkaClient(brokers)
-
- if producer == 'keyed':
- ProducerClass = KeyedProducer
- else:
- ProducerClass = SimpleProducer
-
- kafka_producer = ProducerClass(kafka, **kafka_producer_args)
+ kafka = KafkaClient(hosts=brokers)
# These will be used if incoming events are not interpolatable.
default_topic = topic.encode('utf8')
default_key = key.encode('utf8')
-
- kafka_topic_create_timeout_seconds = 0.1
if blacklist:
blacklist_pattern = re.compile(blacklist)
@@ -177,50 +157,22 @@
continue
message_topic = (topic % event).encode('utf8')
- if producer == 'keyed':
- message_key = (key % event).encode('utf8')
+ message_key = (key % event).encode('utf8')
else:
message_topic = default_topic
message_key = default_key
- try:
- # Make sure this topic exists before we attempt to produce to it.
- # This call will timeout in kafka_topic_create_timeout_seconds.
- # This should return faster than this if this kafka client has
- # already cached topic metadata for this topic. Otherwise
- # it will try to ask Kafka for it each time. Make sure
- # auto.create.topics.enabled is true for your Kafka cluster!
- kafka.ensure_topic_exists(
- message_topic,
- kafka_topic_create_timeout_seconds
- )
- except KafkaTimeoutError:
- error_message = "Failed to ensure Kafka topic %s exists " \
- "in %f seconds when producing event" % (
- message_topic,
- kafka_topic_create_timeout_seconds
- )
- if isinstance(event, dict):
- error_message += " of schema %s revision %d" % (
- event['schema'],
- event['revision']
- )
- error_message += ". Skipping event. " \
- "(This might be ok if this is a new topic.)"
- logging.warn(error_message)
- continue
+ # Get a pykafka Topic instance for the topic
+ kafka_topic = kafka.topics[message_topic]
+ # Get a Producer instance for the topic
+ kafka_producer = kafka_topic.get_producer(**kafka_producer_args)
if raw:
value = event.encode('utf-8')
else:
value = json.dumps(event, sort_keys=True)
- # send_messages() for the different producer types have different
- # signatures. Call it appropriately.
- if producer == 'keyed':
- kafka_producer.send_messages(message_topic, message_key, value)
- else:
- kafka_producer.send_messages(message_topic, value)
+ kafka_producer.produce([(message_key, value)])
def insert_stats(stats, inserted_count):
@@ -461,7 +413,7 @@
if k in inspect.getargspec(BalancedConsumer.__init__).args
}
- kafka_client = PyKafkaClient(hosts=brokers)
+ kafka_client = KafkaClient(hosts=brokers)
kafka_topic = kafka_client.topics[topic]
consumer = kafka_topic.get_balanced_consumer(
--
To view, visit https://gerrit.wikimedia.org/r/232408
To unsubscribe, visit https://gerrit.wikimedia.org/r/settings
Gerrit-MessageType: newchange
Gerrit-Change-Id: I5a1b801316786ba017836946ed6e96258187695b
Gerrit-PatchSet: 1
Gerrit-Project: mediawiki/extensions/EventLogging
Gerrit-Branch: master
Gerrit-Owner: Madhuvishy <[email protected]>
_______________________________________________
MediaWiki-commits mailing list
[email protected]
https://lists.wikimedia.org/mailman/listinfo/mediawiki-commits