Repository navigation
Expand file tree
/
Copy pathtransactional_producer.py
More file actions
62 lines (51 loc) · 1.63 KB
/
Copy pathtransactional_producer.py
File metadata and controls
62 lines (51 loc) · 1.63 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
from confluent_kafka import Producer
from confluent_kafka.admin import AdminClient, NewTopic
from datetime import datetime
from avro import schema, io
from avro.datafile import DataFileWriter
# Kafka broker configuration
bootstrap_servers = 'localhost:9092'
topic = 'Rajeev'
avro_schema_str = """
{
"type": "record",
"name": "Message",
"fields": [
{"name": "id", "type": "int"},
{"name": "content", "type": "string"}
]
}
"""
# Define Avro schema
avro_schema = schema.Parse(avro_schema_str)
# Create a Kafka producer instance
producer = Producer({
'bootstrap.servers': bootstrap_servers,
'transactional.id': 'my_transactional_producer'
})
# Initialize the producer as a transactional producer
producer.init_transactions()
# Start the transaction
producer.begin_transaction()
try:
# Create Avro message
message_value = {"id": 1, "content": "Hello, Avro!"}
# Serialize the Avro message
avro_writer = io.DatumWriter(avro_schema)
avro_bytes_writer = io.BytesIO()
data_file_writer = DataFileWriter(avro_bytes_writer, avro_writer, avro_schema)
data_file_writer.append(message_value)
data_file_writer.flush()
serialized_message = avro_bytes_writer.getvalue()
# Produce the Avro message within the transaction
producer.produce(topic=topic, value=serialized_message)
# Commit the transaction
producer.commit_transaction()
print("Transaction committed successfully.")
except Exception as e:
# Abort the transaction in case of any exception
producer.abort_transaction()
print(f"Transaction aborted: {str(e)}")
finally:
# Close the producer
producer.close()