По итогам тестирования
This commit is contained in:
@@ -17,7 +17,8 @@ GP_CONN_ID=greenplum_conn
|
|||||||
GP_USE_AIRFLOW_CONN=true
|
GP_USE_AIRFLOW_CONN=true
|
||||||
|
|
||||||
# Kafka Configuration
|
# Kafka Configuration
|
||||||
KAFKA_BOOTSTRAP=kafka:9092
|
KAFKA_BOOTSTRAP=kafka:29092
|
||||||
KAFKA_TOPIC=orders
|
KAFKA_TOPIC=orders
|
||||||
KAFKA_BATCH_SIZE=500
|
KAFKA_BATCH_SIZE=500
|
||||||
KAFKA_POLL_TIMEOUT=10
|
KAFKA_POLL_TIMEOUT=10
|
||||||
|
KAFKA_MAX_EMPTY_POLLS=3
|
||||||
|
|||||||
@@ -77,10 +77,12 @@ airflow connections add 'greenplum_conn' \
|
|||||||
- Greenplum: `localhost:${GP_PORT:-5432}` (внешний порт проброшен из контейнера)
|
- Greenplum: `localhost:${GP_PORT:-5432}` (внешний порт проброшен из контейнера)
|
||||||
- Postgres (Airflow metadata): `localhost:5433`
|
- Postgres (Airflow metadata): `localhost:5433`
|
||||||
- Kafka (для клиентов на хосте): `localhost:9092`
|
- Kafka (для клиентов на хосте): `localhost:9092`
|
||||||
|
- Kafka (из контейнеров Docker): `kafka:29092`
|
||||||
|
|
||||||
### Параметры чтения/загрузки
|
### Параметры чтения/загрузки
|
||||||
- `KAFKA_BATCH_SIZE` — размер батча при вставке в Greenplum (по умолчанию 500).
|
- `KAFKA_BATCH_SIZE` — размер батча при вставке в Greenplum (по умолчанию 500).
|
||||||
- `KAFKA_POLL_TIMEOUT` — таймаут ожидания сообщения в секундах (по умолчанию 10).
|
- `KAFKA_POLL_TIMEOUT` — таймаут ожидания сообщения в секундах (по умолчанию 10).
|
||||||
|
- `KAFKA_MAX_EMPTY_POLLS` — сколько подряд пустых `poll` допускается перед выходом из цикла (по умолчанию 3).
|
||||||
|
|
||||||
### Проверка загрузки
|
### Проверка загрузки
|
||||||
```bash
|
```bash
|
||||||
|
|||||||
Binary file not shown.
@@ -13,10 +13,11 @@ from psycopg2.extras import execute_values
|
|||||||
|
|
||||||
from helpers.greenplum import get_gp_conn
|
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")
|
TOPIC = os.getenv("KAFKA_TOPIC", "orders")
|
||||||
BATCH_SIZE = int(os.getenv("KAFKA_BATCH_SIZE", "500"))
|
BATCH_SIZE = int(os.getenv("KAFKA_BATCH_SIZE", "500"))
|
||||||
POLL_TIMEOUT_S = int(os.getenv("KAFKA_POLL_TIMEOUT", "10"))
|
POLL_TIMEOUT_S = int(os.getenv("KAFKA_POLL_TIMEOUT", "10"))
|
||||||
|
MAX_EMPTY_POLLS = int(os.getenv("KAFKA_MAX_EMPTY_POLLS", "3"))
|
||||||
def _create_table():
|
def _create_table():
|
||||||
ddl = """
|
ddl = """
|
||||||
CREATE TABLE IF NOT EXISTS public.orders (
|
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:
|
with get_gp_conn() as conn, conn.cursor() as cur:
|
||||||
batch: List[Tuple] = []
|
batch: List[Tuple] = []
|
||||||
consumed = 0
|
consumed = 0
|
||||||
while consumed < max_messages:
|
empty_polls = 0
|
||||||
|
while consumed < max_messages and empty_polls < MAX_EMPTY_POLLS:
|
||||||
msg = consumer.poll(timeout_s)
|
msg = consumer.poll(timeout_s)
|
||||||
if msg is None:
|
if msg is None:
|
||||||
# Нет новых сообщений — сбрасываем остаток батча и выходим
|
empty_polls += 1
|
||||||
if batch:
|
continue
|
||||||
_flush_batch(cur, batch)
|
empty_polls = 0
|
||||||
conn.commit()
|
|
||||||
batch.clear()
|
|
||||||
break
|
|
||||||
if msg.error():
|
if msg.error():
|
||||||
raise KafkaException(msg.error())
|
raise KafkaException(msg.error())
|
||||||
|
|
||||||
@@ -102,7 +101,7 @@ def _consume_and_load(max_messages=1000, timeout_s: Optional[int] = None):
|
|||||||
batch.append(
|
batch.append(
|
||||||
(
|
(
|
||||||
int(data["order_id"]),
|
int(data["order_id"]),
|
||||||
data["order_ts"],
|
datetime.fromisoformat(data["order_ts"]),
|
||||||
int(data["customer_id"]),
|
int(data["customer_id"]),
|
||||||
float(data["amount"]),
|
float(data["amount"]),
|
||||||
)
|
)
|
||||||
|
|||||||
Reference in New Issue
Block a user