68 lines
1.5 KiB
YAML
68 lines
1.5 KiB
YAML
---
|
|
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",
|
|
]
|