Загрузка данных
"""Загрузка из DWH, кластеризация, запись в core (DAG: delete_full_insert)."""
from __future__ import annotations
from datetime import datetime
from typing import Any
import pandas as pd
from processing_class.mssql_pyodbc_hook import MsSqlPyodbcHook
from service import decorator_func as dec
from sql_script import custom_select_source as custom_select
from processing_class.nomenclature_fuzzy_groups.cluster import cluster_fuzzy, fuzzy_block
from processing_class.nomenclature_fuzzy_groups.constants import (
FUZZY_BLOCK_PREFIX_LEN,
FUZZY_CORE_THRESHOLD,
SINGLE_TOKEN_TAIL_MIN_RATIO,
VARIANT_TAIL_SPLIT_THRESHOLD,
)
from processing_class.nomenclature_fuzzy_groups.naming import normalize_nomenclature, parse_name
def load_nomenclature(hook: Any) -> pd.DataFrame:
return hook.get_pandas_df(sql=custom_select.nomenclature_fuzzy_groups)
def nomen_group_cluster_max(s: pd.Series) -> object:
strs = [str(x).strip() for x in s.dropna() if str(x).strip()]
if not strs:
return pd.NA
num = pd.to_numeric(pd.Series(strs, dtype=object), errors="coerce")
if num.notna().all():
mx = num.max()
tied = [t for t, n in zip(strs, num) if n == mx]
return max(tied)
return max(strs)
def assign_cluster_nomen_group(df: pd.DataFrame) -> pd.DataFrame:
out = df.copy()
out["nomen_group_cluster"] = out.groupby("fuzzy_cluster_id", sort=False)["nomen_group"].transform(
nomen_group_cluster_max
)
return out
def process_dataframe(df: pd.DataFrame) -> pd.DataFrame:
"""Кластеризация по загруженному DataFrame (без SQL)."""
out = df.copy()
out["normalized_name"] = out["nomenclature_name"].map(
lambda x: normalize_nomenclature(x if pd.notna(x) else "")
)
parsed = [parse_name(n) for n in out["normalized_name"].tolist()]
out["name_core"] = [p.core for p in parsed]
packs = [p.pack for p in parsed]
vars_only = [p.variant for p in parsed]
out["tail_pack"] = packs
out["variant_tail"] = [f"{v} {p}".strip() if p else v for v, p in zip(vars_only, packs)]
out["fuzzy_block"] = [fuzzy_block(n, FUZZY_BLOCK_PREFIX_LEN) for n in out["normalized_name"]]
out["fuzzy_cluster_id"] = cluster_fuzzy(
out["name_core"].tolist(),
vars_only,
packs,
out["fuzzy_block"].tolist(),
FUZZY_CORE_THRESHOLD,
VARIANT_TAIL_SPLIT_THRESHOLD,
SINGLE_TOKEN_TAIL_MIN_RATIO,
)
out = assign_cluster_nomen_group(out)
return out.sort_values(
by=["nomen_group_cluster", "nomenclature_code_first"],
ascending=[True, True],
na_position="last",
ignore_index=True,
)
class NomenclatureFuzzyGroups:
"""Расчёт fuzzy-групп и полная перезапись таблицы в core (TRUNCATE + INSERT)."""
def __init__(
self,
connect_id: str,
db_core: str,
schema_core: str,
table_name_core: str,
fields_core: list,
):
self._mssql_hook = MsSqlPyodbcHook(mssql_conn_id=connect_id)
self._fields_core = fields_core
self._table_name_core_path = f"[{db_core}].[{schema_core}].[{table_name_core}]"
def _bulk_insert_core(self, df: pd.DataFrame, commit_every: int = 5000) -> None:
cols = self._fields_core
col_list = ", ".join(f"[{c}]" for c in cols) + ", [dlm$]"
placeholders = ", ".join("?" * len(cols)) + ", ?"
insert_sql = (
f"INSERT INTO {self._table_name_core_path} ({col_list}) "
f"VALUES ({placeholders})"
)
dlm = datetime.now()
rows = [
tuple(None if pd.isna(v) else v for v in row) + (dlm,)
for row in df[cols].itertuples(index=False, name=None)
]
if not rows:
return
conn = self._mssql_hook.get_conn()
try:
cur = conn.cursor()
cur.fast_executemany = True
for start in range(0, len(rows), commit_every):
cur.executemany(insert_sql, rows[start:start + commit_every])
conn.commit()
finally:
conn.close()
@dec.task_python_operator
@dec.timer('Nomenclature fuzzy groups delete_full_insert')
def delete_full_insert(self, **kwargs):
result = process_dataframe(load_nomenclature(self._mssql_hook))[self._fields_core]
self._mssql_hook.get_records(f"TRUNCATE TABLE {self._table_name_core_path}")
self._bulk_insert_core(result)
print(f'Записано в {self._table_name_core_path}: {len(result):,} строк')
--------------------------------------------------------------------------------------------------------------------------------------------------------
from airflow.models.dag import DAG
from processing_class.nomenclature_fuzzy_groups.processing_nomenclature import NomenclatureFuzzyGroups
from service.notify_failure_cooldown import notify_failure_cooldown
import pendulum
default_args = {
'owner': 'nomenclature',
'depends_on_past': False,
'retries': 0,
'on_failure_callback': notify_failure_cooldown(60),
'start_date': pendulum.datetime(2026, 4, 1, tz='Europe/Moscow'),
}
connector_id = 'mssql_dwh'
db_core = 'dwh'
schema_core = 'dbo'
with DAG(
dag_id='nomenclature_fuzzy_groups',
default_args=default_args,
tags=['nomenclature', 'fuzzy', 'dwh'],
schedule_interval='0 20 * * *',
catchup=False
):
table_name_core = 'nomenclature_fuzzy_groups'
fields_core = [
'nomenclature_name', 'normalized_name', 'nomen_group', 'nomen_group_cluster',
'group_buyer_cnt', 'group_rebrand_cnt', 'nomenclature_code_customer_first',
'rebrand_sku_erp_first', 'nomenclature_code_first',
]
main = NomenclatureFuzzyGroups(
connector_id, db_core, schema_core, table_name_core, fields_core,
)
main.delete_full_insert()