From 4f8d9fa9dcbffb961b546465e51326dafeb2e6c1 Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Fri, 18 Mar 2016 11:25:59 -0700 Subject: Add producer.flush() to usage docs --- docs/index.rst | 8 ++++++-- docs/usage.rst | 7 +++++++ 2 files changed, 13 insertions(+), 2 deletions(-) (limited to 'docs') diff --git a/docs/index.rst b/docs/index.rst index d8f826a..eb8f429 100644 --- a/docs/index.rst +++ b/docs/index.rst @@ -74,9 +74,13 @@ client. See `KafkaProducer `_ for more details. >>> from kafka import KafkaProducer >>> producer = KafkaProducer(bootstrap_servers='localhost:1234') ->>> producer.send('foobar', b'some_message_bytes') +>>> for _ in range(100): +... producer.send('foobar', b'some_message_bytes') ->>> # Blocking send +>>> # Block until all pending messages are sent +>>> producer.flush() + +>>> # Block until a single message is sent (or timeout) >>> producer.send('foobar', b'another_message').get(timeout=60) >>> # Use a key for hashed-partitioning diff --git a/docs/usage.rst b/docs/usage.rst index d48cc0a..85fc44f 100644 --- a/docs/usage.rst +++ b/docs/usage.rst @@ -87,5 +87,12 @@ KafkaProducer producer = KafkaProducer(value_serializer=lambda m: json.dumps(m).encode('ascii')) producer.send('json-topic', {'key': 'value'}) + # produce asynchronously + for _ in range(100): + producer.send('my-topic', b'msg') + + # block until all async messages are sent + producer.flush() + # configure multiple retries producer = KafkaProducer(retries=5) -- cgit v1.2.1