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

Reply via email to