This commit is contained in:
ivan 2026-08-31 14:49:47 +05:00
parent ba72d321ba
commit 96e42f9899
7 changed files with 141 additions and 3 deletions

View File

@ -86,7 +86,32 @@ spec:
labels:
monitoring: prometheus
volumes:
_default:
- name: kafka-configmap
mountPath:
_default: /opt/src/flow/kafka.py
subPath:
_default: kafka.py
readOnly:
_default: true
configMap:
name:
_default: kafka-configmap
items:
- key: kafka.py
path:
_default: kafka.py
envs:
- name: KAFKA_HOST
value:
_default: "brusnika-stage-kafka-bootstrap.kafka.svc.cluster.local:9093"
- name: KAFKA_SSL_CERT
value:
_default: ""
- name: LOG_LEVEL
value:
_default: "DEBUG"
@ -256,6 +281,16 @@ spec:
_default: "rabbitmq-secret"
secretKey: "vhost"
- name: KAFKA_USERNAME
secretName:
_default: "kafka-secret"
secretKey: "username"
- name: KAFKA_PASSWORD
secretName:
_default: "kafka-secret"
secretKey: "password"
commitSha: ""
gitlabUri: ""
gitlabJobUrl: ""

View File

@ -86,7 +86,32 @@ spec:
labels:
monitoring: prometheus
volumes:
_default:
- name: kafka-configmap
mountPath:
_default: /opt/src/flow/kafka.py
subPath:
_default: kafka.py
readOnly:
_default: true
configMap:
name:
_default: kafka-configmap
items:
- key: kafka.py
path:
_default: kafka.py
envs:
- name: KAFKA_HOST
value:
_default: "brusnika-stage-kafka-bootstrap.kafka.svc.cluster.local:9093"
- name: KAFKA_SSL_CERT
value:
_default: ""
- name: LOG_LEVEL
value:
_default: "DEBUG"
@ -254,6 +279,16 @@ spec:
_default: "rabbitmq-secret"
secretKey: "vhost"
- name: KAFKA_USERNAME
secretName:
_default: "kafka-secret"
secretKey: "username"
- name: KAFKA_PASSWORD
secretName:
_default: "kafka-secret"
secretKey: "password"
commitSha: ""
gitlabUri: ""
gitlabJobUrl: ""

View File

@ -0,0 +1,67 @@
---
apiVersion: v1
kind: ConfigMap
metadata:
name: kafka-configmap
namespace: flows
data:
kafka.py: |
import ssl
from functools import partial
import anyio
from faststream.kafka import KafkaBroker
from faststream.security import BaseSecurity
from flow.config import logger, settings
from flow.models.events import BaseEvent
def _build_broker() -> KafkaBroker:
ssl_context = ssl.create_default_context()
ssl_context.check_hostname = False
ssl_context.verify_mode = ssl.CERT_NONE
if settings.kafka.ssl_cert:
ssl_context.load_verify_locations(cadata=settings.kafka.ssl_cert)
security = BaseSecurity(
ssl_context=ssl_context,
use_ssl=True,
)
return KafkaBroker(
bootstrap_servers=[settings.kafka.host],
security=security,
)
broker = _build_broker()
async def _publish(event: BaseEvent, key: str) -> None:
await broker.publish(
message=event,
topic=event.metadata.event_type,
key=key.encode(),
)
def publish_event(event: BaseEvent, key: str) -> None:
if not settings.enable_events:
return
try:
anyio.from_thread.run(partial(_publish, event, key))
except Exception:
logger.exception(
"Failed to publish kafka event %s",
event.metadata.event_type,
)
__all__ = [
"broker",
"publish_event",
]

View File

@ -3,6 +3,7 @@ apiVersion: kustomize.config.k8s.io/v1beta1
kind: Kustomization
namespace: flows
resources:
- kafka-configmap.yaml
- frontend.yaml
- backend.yaml
- celery.yaml

View File

@ -141,7 +141,7 @@ spec:
- name: KAFKA_BROKERS
value:
_default: "local:9091"
_default: "brusnika-stage-kafka-bootstrap.kafka.svc.cluster.local:9093"
- name: KAFKA_SSL_CAFILE
value:

View File

@ -144,7 +144,7 @@ spec:
- name: KAFKA_HOST
value:
_default: "donstroi-kafka-bootstrap.kafka.svc.cluster.local"
_default: "brusnika-stage-kafka-bootstrap.kafka.svc.cluster.local"
- name: KAFKA_PORT
value:

View File

@ -93,7 +93,7 @@ spec:
- name: KAFKA_BROKERS
value:
_default: "rc1d-s0a1ujcbj6fdk26b.mdb.yandexcloud.net:9091"
_default: "brusnika-stage-kafka-bootstrap.kafka.svc.cluster.local:9093"
- name: KAFKA_GROUP
value: