Skip to content

Commit b1efd86

Browse files
authored
fix kafka admin error (#16)
1 parent 267d3c2 commit b1efd86

File tree

2 files changed

+11
-1
lines changed

2 files changed

+11
-1
lines changed

backend/deepchecks_monitoring/bgtasks/model_version_topic_delete.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -48,7 +48,7 @@ async def run(self, task: 'Task', session: AsyncSession, resources_provider: Res
4848
if self.kafka_admin is None:
4949
with self.lock:
5050
if self.kafka_admin is None:
51-
self.kafka_admin = AIOKafkaAdminClient(**resources_provider.kafka_settings.kafka_params)
51+
self.kafka_admin = AIOKafkaAdminClient(**resources_provider.kafka_settings.kafka_admin_params)
5252
await self.kafka_admin.start()
5353

5454
# Backward compatibility, remove in next release and replace with:

backend/deepchecks_monitoring/config.py

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -79,6 +79,16 @@ def kafka_params(self):
7979
'metadata_max_age_ms': self.kafka_max_metadata_age
8080
}
8181

82+
@property
83+
def kafka_admin_params(self):
84+
"""Get connection parameters for kafka admin."""
85+
return {
86+
'bootstrap_servers': self.kafka_host,
87+
'security_protocol': self.kafka_security_protocol,
88+
'ssl_context': create_ssl_context(),
89+
'metadata_max_age_ms': self.kafka_max_metadata_age
90+
}
91+
8292

8393
class DatabaseSettings(BaseDeepchecksSettings):
8494
"""Database settings."""

0 commit comments

Comments
 (0)