diff --git a/airflow-greenplum/.env.example b/airflow-greenplum/.env.example index fb6257a..7d5f63a 100644 --- a/airflow-greenplum/.env.example +++ b/airflow-greenplum/.env.example @@ -17,7 +17,8 @@ GP_CONN_ID=greenplum_conn GP_USE_AIRFLOW_CONN=true # Kafka Configuration -KAFKA_BOOTSTRAP=kafka:9092 +KAFKA_BOOTSTRAP=kafka:29092 KAFKA_TOPIC=orders KAFKA_BATCH_SIZE=500 KAFKA_POLL_TIMEOUT=10 +KAFKA_MAX_EMPTY_POLLS=3 diff --git a/airflow-greenplum/README.md b/airflow-greenplum/README.md index 8311054..e9a4280 100644 --- a/airflow-greenplum/README.md +++ b/airflow-greenplum/README.md @@ -77,10 +77,12 @@ airflow connections add 'greenplum_conn' \ - Greenplum: `localhost:${GP_PORT:-5432}` (внешний порт проброшен из контейнера) - Postgres (Airflow metadata): `localhost:5433` - Kafka (для клиентов на хосте): `localhost:9092` +- Kafka (из контейнеров Docker): `kafka:29092` ### Параметры чтения/загрузки - `KAFKA_BATCH_SIZE` — размер батча при вставке в Greenplum (по умолчанию 500). - `KAFKA_POLL_TIMEOUT` — таймаут ожидания сообщения в секундах (по умолчанию 10). +- `KAFKA_MAX_EMPTY_POLLS` — сколько подряд пустых `poll` допускается перед выходом из цикла (по умолчанию 3). ### Проверка загрузки ```bash diff --git a/airflow-greenplum/airflow/dags/helpers/__pycache__/greenplum.cpython-312.pyc b/airflow-greenplum/airflow/dags/helpers/__pycache__/greenplum.cpython-312.pyc new file mode 100644 index 0000000..5098ff6 Binary files /dev/null and b/airflow-greenplum/airflow/dags/helpers/__pycache__/greenplum.cpython-312.pyc differ diff --git a/airflow-greenplum/airflow/dags/kafka_to_greenplum.py b/airflow-greenplum/airflow/dags/kafka_to_greenplum.py index a3e92f0..969f2b9 100644 --- a/airflow-greenplum/airflow/dags/kafka_to_greenplum.py +++ b/airflow-greenplum/airflow/dags/kafka_to_greenplum.py @@ -13,10 +13,11 @@ from psycopg2.extras import execute_values from helpers.greenplum import get_gp_conn -KAFKA_BOOTSTRAP = os.getenv("KAFKA_BOOTSTRAP", "kafka:9092") +KAFKA_BOOTSTRAP = os.getenv("KAFKA_BOOTSTRAP", "kafka:29092") TOPIC = os.getenv("KAFKA_TOPIC", "orders") BATCH_SIZE = int(os.getenv("KAFKA_BATCH_SIZE", "500")) POLL_TIMEOUT_S = int(os.getenv("KAFKA_POLL_TIMEOUT", "10")) +MAX_EMPTY_POLLS = int(os.getenv("KAFKA_MAX_EMPTY_POLLS", "3")) def _create_table(): ddl = """ CREATE TABLE IF NOT EXISTS public.orders ( @@ -86,15 +87,13 @@ def _consume_and_load(max_messages=1000, timeout_s: Optional[int] = None): with get_gp_conn() as conn, conn.cursor() as cur: batch: List[Tuple] = [] consumed = 0 - while consumed < max_messages: + empty_polls = 0 + while consumed < max_messages and empty_polls < MAX_EMPTY_POLLS: msg = consumer.poll(timeout_s) if msg is None: - # Нет новых сообщений — сбрасываем остаток батча и выходим - if batch: - _flush_batch(cur, batch) - conn.commit() - batch.clear() - break + empty_polls += 1 + continue + empty_polls = 0 if msg.error(): raise KafkaException(msg.error()) @@ -102,7 +101,7 @@ def _consume_and_load(max_messages=1000, timeout_s: Optional[int] = None): batch.append( ( int(data["order_id"]), - data["order_ts"], + datetime.fromisoformat(data["order_ts"]), int(data["customer_id"]), float(data["amount"]), )