This commit is contained in:
CHIEFSOFT\ameye
2025-10-17 18:24:07 -04:00
parent f87953114c
commit 4bbb4ca17b
2 changed files with 53 additions and 0 deletions
+16
View File
@@ -0,0 +1,16 @@
# Flask and Extensions
Flask==2.3.3
Flask-Marshmallow==0.15.0
marshmallow==3.19.0
Flask-Cors==3.0.10
gunicorn
requests
confluent-kafka==1.9.2
flask-sqlalchemy
psycopg2-binary
alembic
python-dateutil
oracledb
Flask-Mail==0.10.0
pandas==2.1.3
openpyxl==3.1.5
+37
View File
@@ -0,0 +1,37 @@
import threading
from app import create_app
from app.integrations import KafkaIntegration
from app.config import settings
from app.utils.logger import logger
app = create_app()
kafka = KafkaIntegration()
def start_kafka_consumer(app):
with app.app_context():
logger.info("Starting Kafka consumer...")
while True:
try:
message = kafka.receive_messages(
topics=settings.KAFKA_TOPICS, timeout=settings.KAFKA_TIMEOUT
)
if message:
logger.info(f"Processed message: {message}")
else:
logger.info("No message received within timeout")
except Exception as e:
logger.error(f"Error while receiving message: {e}")
if __name__ != "__main__":
# Expose WSGI app instance for Gunicorn
wsgi_app = app
# Start kafka in a thread
threading.Thread(target=start_kafka_consumer, args=(app,), daemon=True).start()