(20250107) Added a few params to synchronous Kafka.
This commit is contained in:
@@ -138,7 +138,10 @@ class ProducerKafka:
|
|||||||
ca_file: str | None = None,
|
ca_file: str | None = None,
|
||||||
cert_file: str | None = None,
|
cert_file: str | None = None,
|
||||||
key_file: str | None = None,
|
key_file: str | None = None,
|
||||||
client_id: str | int | None = None
|
client_id: str | int | None = None,
|
||||||
|
acks: int = 0,
|
||||||
|
retries: int = 1,
|
||||||
|
linger_ms: int = 0
|
||||||
):
|
):
|
||||||
|
|
||||||
"""
|
"""
|
||||||
@@ -156,7 +159,10 @@ class ProducerKafka:
|
|||||||
if not isinstance(bootstrap_servers, list): bootstrap_servers = [bootstrap_servers]
|
if not isinstance(bootstrap_servers, list): bootstrap_servers = [bootstrap_servers]
|
||||||
config = {
|
config = {
|
||||||
"bootstrap.servers": ",".join(bootstrap_servers),
|
"bootstrap.servers": ",".join(bootstrap_servers),
|
||||||
"security.protocol": security_protocol
|
"security.protocol": security_protocol,
|
||||||
|
"acks": acks,
|
||||||
|
"retries": retries,
|
||||||
|
"linger.ms": linger_ms
|
||||||
}
|
}
|
||||||
|
|
||||||
# Add the SSL security details:
|
# Add the SSL security details:
|
||||||
|
|||||||
Reference in New Issue
Block a user