from fansifter_common.utils.functional import lazy_proxy from app.adapters.kafka import kafka_producer from app.config import settings from .base import DataSink from .dummy import DummySink from .kafka import KafkaSink from .postgres import PostgresSink def get_sink() -> DataSink: if settings.data_sink_backend == "kafka": return KafkaSink( producer=kafka_producer, topic_fans=settings.kafka_topic_fans, topic_top_artists=settings.kafka_topic_top_artists, topic_top_tracks=settings.kafka_topic_top_tracks, topic_recently_played=settings.kafka_topic_recently_played, topic_playlists=settings.kafka_topic_playlists, topic_saved_albums=settings.kafka_topic_saved_albums, topic_saved_tracks=settings.kafka_topic_saved_tracks, topic_followed_artists=settings.kafka_topic_followed_artists, ) if settings.data_sink_backend == "dummy": return DummySink() return PostgresSink() data_sink = lazy_proxy(get_sink)