@@ -66,14 +66,14 @@ pub fn spawn_kafka_task(
6666 } ;
6767
6868 // Initialize transactional producer if configured.
69- if config. transactional {
70- if let Err ( e) = producer. init_transactions ( Duration :: from_secs ( 10 ) ) {
71- warn ! (
72- stream = %stream_name ,
73- error = %e ,
74- "failed to init Kafka transactions — falling back to at-least-once"
75- ) ;
76- }
69+ if config. transactional
70+ && let Err ( e) = producer. init_transactions ( Duration :: from_secs ( 10 ) )
71+ {
72+ warn ! (
73+ stream = %stream_name ,
74+ error = %e ,
75+ "failed to init Kafka transactions — falling back to at-least-once"
76+ ) ;
7777 }
7878
7979 let group_name = format ! ( "_kafka_{stream_name}" ) ;
@@ -99,11 +99,11 @@ pub fn spawn_kafka_task(
9999 let batch_size = events. events. len( ) ;
100100
101101 // Begin transaction if configured.
102- if config. transactional {
103- if let Err ( e) = producer. begin_transaction( ) {
104- warn! ( error = %e , "Kafka begin_transaction failed" ) ;
105- continue ;
106- }
102+ if config. transactional
103+ && let Err ( e) = producer. begin_transaction( )
104+ {
105+ warn! ( error = %e , "Kafka begin_transaction failed" ) ;
106+ continue ;
107107 }
108108
109109 let mut published = 0u32 ;
@@ -136,13 +136,12 @@ pub fn spawn_kafka_task(
136136 }
137137
138138 // Commit Kafka transaction.
139- if config. transactional && published > 0 {
140- if let Err ( e) = producer
139+ if config. transactional && published > 0
140+ && let Err ( e) = producer
141141 . commit_transaction( Duration :: from_secs( 10 ) )
142- {
143- warn!( error = %e, "Kafka commit_transaction failed" ) ;
144- continue ;
145- }
142+ {
143+ warn!( error = %e, "Kafka commit_transaction failed" ) ;
144+ continue ;
146145 }
147146
148147 // Commit consumer offsets for successfully published events.
@@ -189,7 +188,7 @@ fn create_producer(
189188
190189 if config. transactional {
191190 client_config. set ( "enable.idempotence" , "true" ) ;
192- client_config. set ( "transactional.id" , & format ! ( "nodedb-kafka-{stream_name}" ) ) ;
191+ client_config. set ( "transactional.id" , format ! ( "nodedb-kafka-{stream_name}" ) ) ;
193192 }
194193
195194 client_config
0 commit comments