
broker = KafkaBroker("localhost:9092")
app = FastStream(broker)
@broker.subscriber("raw_texts_topic", group_id="text_parsers")
async def handle_text(text_data: dict):
# Ανάλυση και αποθήκευση στη βάση δεδομένων διαρκεί περίπου 10 δευτερόλεπτα
parse_text_and_save_to_postgres(text_data)