forked from confluentinc/confluent-kafka-python
-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathshare_consumer_avro.py
More file actions
161 lines (134 loc) · 5.86 KB
/
Copy pathshare_consumer_avro.py
File metadata and controls
161 lines (134 loc) · 5.86 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
#!/usr/bin/env python
# -*- coding: utf-8 -*-
#
# Copyright 2026 Confluent Inc.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
# Example KIP-932 DeserializingShareConsumer reading Avro via Schema Registry.
#
# The key/value deserializers go in the config dict and run inside poll(). A
# record that can't be deserialized isn't raised — it comes back with its raw
# bytes and a _KEY/_VALUE_DESERIALIZATION error, so you can REJECT it and move
# on.
import argparse
import sys
from confluent_kafka import (
AcknowledgeType,
ConcurrentModificationException,
DeserializingShareConsumer,
IllegalStateException,
KafkaError,
KafkaException,
)
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroDeserializer
from confluent_kafka.serialization import StringDeserializer
class User(object):
"""
User record
Args:
name (str): User's name
favorite_number (int): User's favorite number
favorite_color (str): User's favorite color
"""
def __init__(self, name=None, favorite_number=None, favorite_color=None):
self.name = name
self.favorite_number = favorite_number
self.favorite_color = favorite_color
def dict_to_user(obj, ctx):
"""
Converts an Avro-decoded object literal (dict) into a User instance.
Args:
obj (dict): Object literal(dict)
ctx (SerializationContext): Metadata pertaining to the serialization
operation.
"""
if obj is None:
return None
return User(name=obj['name'], favorite_number=obj['favorite_number'], favorite_color=obj['favorite_color'])
def main(args):
schema_str = """
{
"namespace": "confluent.io.examples.serialization.avro",
"name": "User",
"type": "record",
"fields": [
{"name": "name", "type": "string"},
{"name": "favorite_number", "type": "long"},
{"name": "favorite_color", "type": "string"}
]
}
"""
schema_registry_client = SchemaRegistryClient({'url': args.schema_registry})
avro_deserializer = AvroDeserializer(schema_registry_client, schema_str, dict_to_user)
# Deserializers go in the config dict.
conf = {
'bootstrap.servers': args.bootstrap_servers,
'group.id': args.group,
'share.acknowledgement.mode': 'explicit',
'key.deserializer': StringDeserializer('utf_8'),
'value.deserializer': avro_deserializer,
}
sc = DeserializingShareConsumer(conf)
sc.subscribe([args.topic])
try:
while True:
try:
messages = sc.poll(timeout=1.0) # a list, possibly empty
for msg in messages:
err = msg.error()
if err is not None:
if err.code() in (KafkaError._KEY_DESERIALIZATION, KafkaError._VALUE_DESERIALIZATION):
# A record we received but can't decode. In explicit
# mode we still have to ack it — RELEASE it so that
# other consumer can pick and redeliver for processing
# again. In implicit ack mode, we currently don't release
# the record in the Preview internally.
sc.acknowledge(msg, AcknowledgeType.RELEASE)
else:
# Any other flagged record is acked internally by the
# library — see share_consumer.py for the error handling.
sys.stderr.write('%% Error: %s\n' % err)
continue
# value is already deserialized.
user = msg.value()
if user is not None:
print(
"User record {}: name: {}, favorite_number: {}, favorite_color: {}".format(
msg.key(), user.name, user.favorite_number, user.favorite_color
)
)
sc.acknowledge(msg, AcknowledgeType.ACCEPT)
except KafkaException as e:
# Re-raise fatal errors; otherwise log and keep going.
if e.args[0].fatal():
raise
sys.stderr.write('%% Consumer error: %s\n' % e)
continue
except (IllegalStateException, ConcurrentModificationException) as e:
# These signal misuse (polling or acking when not subscribed/closed,
# or from more than one thread), not a transient hiccup — no point
# looping, so bail out.
sys.stderr.write('%% Fatal: %s\n' % e)
raise
except KeyboardInterrupt:
sys.stderr.write('%% Aborted by user\n')
finally:
sc.close()
if __name__ == '__main__':
parser = argparse.ArgumentParser(description="DeserializingShareConsumer Avro example")
parser.add_argument('-b', dest="bootstrap_servers", required=True, help="Bootstrap broker(s) (host[:port])")
parser.add_argument('-s', dest="schema_registry", required=True, help="Schema Registry (http(s)://host[:port])")
parser.add_argument('-t', dest="topic", default="example_serde_avro", help="Topic name")
parser.add_argument('-g', dest="group", default="example_share_serde_avro", help="Share group")
main(parser.parse_args())