Skip to content

feat: PostgreSQL sink (psycopg2 COPY FROM STDIN) - #4

Merged
starychenko merged 13 commits into
mainfrom
feature/postgresql-sink
Mar 10, 2026
Merged

feat: PostgreSQL sink (psycopg2 COPY FROM STDIN)#4
starychenko merged 13 commits into
mainfrom
feature/postgresql-sink

Conversation

@starychenko

Copy link
Copy Markdown
Owner

Summary

  • Додано PostgreSQLSink як третій аналітичний sink поряд із ClickHouse та DuckDB
  • Bulk-завантаження через psycopg2 COPY FROM STDIN — найшвидший метод для PostgreSQL
  • Підтримка SSL (sslmode=require) для self-signed сертифікатів

Зміни

  • olap_tool/config.pyPostgreSQLConfig dataclass, load_postgres_from_env(), інтеграція в AppConfig та build_config()
  • olap_tool/sinks.pyPostgreSQLSink: CREATE TABLE IF NOT EXISTS, ALTER TABLE ADD COLUMN IF NOT EXISTS, DELETE (ідемпотентність), COPY FROM STDIN CSV; thread-safe _schema_lock; rollback guard у delete_period/insert
  • olap_tool/queries.pych_onlysink_only, додано PG/POSTGRESQL до всіх перевірок формату
  • olap_tool/runner.py — info-блок підключення + побудова sink
  • olap_tool/cli.py--format приймає pg, postgresql
  • requirements.txtpsycopg2-binary>=2.9.0
  • .env.example — документація PG_* змінних

Конфігурація (.env)

PG_ENABLED=true
PG_HOST=db.lwhs.xyz
PG_PORT=54321
PG_DATABASE=analytics
PG_USER=analytics
PG_PASSWORD=...
PG_SSL_MODE=require

Test Plan

  • python -c "from olap_tool.sinks import PostgreSQLSink; print('OK')" → OK
  • python olap.py --format pg --last-weeks 1 — завантажує дані в PostgreSQL без XLSX
  • python olap.py --last-weeks 1 з PG_ENABLED=true — паралельно з XLSX
  • Перевірити таблицю в PostgreSQL: SELECT COUNT(*) FROM public.sales

🤖 Generated with Claude Code

starychenko and others added 11 commits March 10, 2026 08:57
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
…etup()

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
@chatgpt-codex-connector

Copy link
Copy Markdown

You have reached your Codex usage limits for code reviews. You can see your limits in the Codex usage dashboard.

@claude

claude Bot commented Mar 10, 2026

Copy link
Copy Markdown

Claude finished @starychenko's task in 3m 13s —— View job


Огляд PR #4: PostgreSQL sink (psycopg2 COPY FROM STDIN)

  • Читаю змінені файли
  • Перевіряю архітектурні контракти (AnalyticsSink, config flow)
  • Перевіряю управління ресурсами та thread safety
  • Перевіряю коректність та data integrity
  • Формую фінальний вердикт

Резюме

PR додає PostgreSQLSink — третій аналітичний sink з bulk-завантаженням через psycopg2 COPY FROM STDIN. Контракт AnalyticsSink виконано повністю, config flow через AppConfig дотримано, ідемпотентність (DELETE WHERE year_num+week_num перед COPY) присутня.

Вердикт: 🔴 Request Changes — знайдено 1 блокер та 3 важливі проблеми.


🔴 BLOCKER

setup() не робить rollback при помилці DDL (olap_tool/sinks.py, рядки 505–509)

with conn.cursor() as cur:
    cur.execute(
        f"CREATE TABLE IF NOT EXISTS {self._full_table()} ({cols_ddl})"
    )
conn.commit()   # ← якщо cur.execute() кинув виняток — сюди не дійде

Якщо CREATE TABLE падає (схема public не існує, недостатньо прав, DDL-конфлікт), psycopg2 залишає з'єднання у стані InFailedSqlTransaction. _get_conn() перевіряє лише self._conn.closed (який буде False — з'єднання технічно відкрите, але зламане). Всі подальші виклики _refresh_schema(), delete_period(), insert() на тому самому з'єднанні кидатимуть psycopg2.errors.InFailedSqlTransaction, які тихо поглинаються в queries.py через except Exception: print_error(...) → дані не потрапляють у PostgreSQL, але виконання продовжується без rollback і без явної діагностики.

Аналогічна проблема є в блоці ALTER TABLE (рядки 519–524) — там є conn.rollback(), але в самому CREATE TABLE блоці його немає.

Виправлення: загорнути весь setup() у try/except з conn.rollback() при помилці. Fix this →


🟡 IMPORTANT

1. Порожній рядок ↔ NULL — втрата розрізнення (olap_tool/sinks.py, рядки 569–576)

buf = io.StringIO()
df.to_csv(buf, index=False, header=False, na_rep="")  # pandas NaN → ""
...
f"FROM STDIN WITH (FORMAT CSV, NULL '')"              # "" → NULL у PostgreSQL

na_rep="" серіалізує NaN як порожній рядок, а NULL '' каже PostgreSQL вважати порожній рядок NULL. Але реальні порожні рядки у текстових колонках також стануть NULL. Втрачається семантична різниця між "" та NULL. Для торговельних даних це може спотворити агрегати (COUNT, COALESCE тощо). Fix this →


2. Подвійний виклик sanitize_df() (olap_tool/queries.py рядок 346, olap_tool/sinks.py рядок 556)

queries.py викликає sanitize_df(df) перед передачею df_for_sinks у всі sink-и, і PostgreSQLSink.insert() викликає sanitize_df(df) ще раз. ClickHouseSink та DuckDBSink.insert() також так роблять — це патерн, що повторюється. Безпечно (функція ідемпотентна), але для великих DataFrame (багато float-колонок) двічі виконується replace([inf, -inf], np.nan). Fix this →


3. _get_conn() не є thread-safe (olap_tool/sinks.py, рядки 458–471)

def _get_conn(self):
    if self._conn is None or self._conn.closed:  # ← race condition без lock
        self._conn = psycopg2.connect(...)
    return self._conn

Якщо два потоки одночасно перевірять self._conn is None, обидва створять з'єднання. Перший _schema_lock захищає схему, але не з'єднання. При поточному використанні (однопотоковий runner) це не проблема, але PostgreSQLSink не документує це обмеження — і додавання batch-скрипту (за аналогією з import_xlsx_to_clickhouse.py) без цього знання призведе до InterfaceError. Fix this →


🔵 NIT

_pandas_dtype_to_pg() розміщена між функціями DuckDB (olap_tool/sinks.py, рядки 157–170)

Функція визначена поміж _pandas_dtype_to_duck() та _to_excel_serial() — тобто всередині DuckDB-блоку, хоча належить до PostgreSQL. Краще перемістити її ближче до PostgreSQLSink.


Що зроблено добре

  • Ідемпотентність реалізована правильно: перевіряє наявність year_num/week_num у схемі перед DELETE, щоб не видалити весь рік при першому запуску (sinks.py, рядки 537–538).
  • Config flow через load_postgres_from_env() у config.py відповідає патерну ClickHouse/DuckDB — os.getenv() тільки в config.py, не розсіяний по модулях.

@claude

claude Bot commented Mar 10, 2026

Copy link
Copy Markdown

Claude finished @starychenko's task in 1m 26s —— View job


Огляд PR #4 (повторна перевірка після оновлення гілки)

  • Читаю змінені файли
  • Перевіряю, чи виправлено попередні зауваження
  • Формую фінальний вердикт

Резюме

Після синхронізації гілки додано 2 фіксуючі коміти (8c90337, 72aa4bf). Блокер виправлено. Дві з трьох "important" проблем закрито (повністю або частково). Залишається одна невирішена — розрізнення "" та NULL.

Вердикт: 🟡 Comment — блокер знято, але 1 важлива проблема залишається відкритою.


Що виправлено ✅

BLOCKER → виправлено (setup(), рядки 508–516)

try:
    with conn.cursor() as cur:
        cur.execute(f"CREATE TABLE IF NOT EXISTS ...")
    conn.commit()
except Exception:
    conn.rollback()
    raise

CREATE TABLE тепер загорнутий у try/except з rollback — відповідає патерну delete_period/insert.

Thread safety → задокументовано (docstring класу, рядки 446–453)

NOT thread-safe: psycopg2-з'єднання не підтримують спільне використання між
потоками. Для batch-скриптів з threading створюйте окремий екземпляр на кожен потік.

Обмеження явно задокументовано — прийнятне рішення.

sanitize_df у setup() — прибрано зайвий виклик із setup() (коміт 72aa4bf).


🟡 Залишається відкритим

Порожній рядок ↔ NULL — втрата розрізнення (olap_tool/sinks.py, рядки 577–583)

df.to_csv(buf, index=False, header=False, na_rep="")   # NaN → ""
# ...
f"FROM STDIN WITH (FORMAT CSV, NULL '')"               # "" → NULL у PostgreSQL

Реальні порожні рядки у текстових колонках ("") стануть NULL у PostgreSQL — семантична різниця втрачається. Для торгових даних це може спотворити COUNT, COALESCE, NOT NULL-умови.

Рекомендоване виправлення — використати sentinel-значення \N (стандарт PostgreSQL для NULL у CSV):

df.to_csv(buf, index=False, header=False, na_rep=r"\N")
# ...
f"FROM STDIN WITH (FORMAT CSV, NULL '\\N')"

Fix this →


🔵 NIT (залишається)

_pandas_dtype_to_pg() розміщена між DuckDB-функціями (sinks.py, рядки 157–170)

Функція визначена між _pandas_dtype_to_duck() та _to_excel_serial(), але належить до PostgreSQL-блоку. Краще перемістити її до секції PostgreSQL (# PostgreSQL sink).


Що зроблено добре

  • Ідемпотентність: DELETE WHERE year_num+week_num присутній, перевіряє наявність колонок у схемі перед DELETE.
  • double sanitize_df: у insert() виклик залишився (як і в DuckDBSink) — це безпечно (ідемпотентно), але варто прибрати для узгодженості після закриття основного issue з NULL ''.
  • Config flow через load_postgres_from_env() у config.pyos.getenv() не виходить за межі config.py.

@claude

claude Bot commented Mar 10, 2026

Copy link
Copy Markdown

Claude finished @starychenko's task in 1m 23s —— View job


Огляд PR #4 (оновлення гілки)

  • Читаю змінені файли
  • Перевіряю, чи виправлено попередні зауваження
  • Формую фінальний вердикт

Результат перевірки

Всі 4 зауваження з попереднього огляду (1 blocker + 3 important) — виправлено. Деталі нижче.

Вердикт: ✅ Approve — блокерів немає, важливих проблем немає.


Що виправлено

🔴 BLOCKER → ✅ Виправлено

setup() тепер загорнуто в try/except з conn.rollback() при помилці DDL (sinks.py, рядки 508–516). Патерн відповідає delete_period() та insert().


🟡 IMPORTANT #1 → ✅ Виправлено

NULL-sentinel \N замість порожнього рядка (sinks.py, рядок 576):

df.to_csv(buf, index=False, header=False, na_rep="\\N")
# ...
r"FROM STDIN WITH (FORMAT CSV, NULL '\N')"

Тепер реальні порожні рядки "" і NULL розрізняються — \N є стандартним sentinel для COPY команд PostgreSQL.


🟡 IMPORTANT #2 → ✅ Виправлено

PostgreSQLSink.insert() більше не викликає sanitize_df() — подвійного виклику немає. queries.py:346 залишається єдиним місцем виклику перед передачею до всіх sink-ів.

Примітка: DuckDBSink.insert():402 ще досі містить власний виклик sanitize_df() — це pre-existing проблема, яку цей PR не вводить і не виправляє. Можна прибрати окремо.


🟡 IMPORTANT #3 → ✅ Виправлено

Клас PostgreSQLSink тепер містить явну документацію thread-safety (sinks.py, рядки 451–452):

NOT thread-safe: psycopg2-з'єднання не підтримують спільне використання між
потоками. Для batch-скриптів з threading створюйте окремий екземпляр на кожен потік.

🔵 NIT → ✅ Виправлено

_pandas_dtype_to_pg() переміщено безпосередньо перед PostgreSQLSink (sinks.py, рядки 429–442), а не в DuckDB-блок.


Що зроблено добре (залишається)

  • Ідемпотентність реалізована правильно: перевіряє наявність year_num/week_num у схемі перед DELETE.
  • Config flow через load_postgres_from_env() дотримано — os.getenv() тільки в config.py.

@starychenko
starychenko merged commit 87049a8 into main Mar 10, 2026
1 check passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant