Files
AWatch-rus/clickhouse-1c/etl/load_1c_exports.py

302 lines
12 KiB
Python

#!/usr/bin/env python3
from __future__ import annotations
import argparse
import csv
import json
import shutil
from dataclasses import dataclass
from datetime import UTC, datetime
from pathlib import Path
from typing import Any
import yaml
from dateutil import parser as date_parser
RAW_TABLES = {
"documents": "raw_1c_documents",
"postings": "raw_1c_postings",
"business_events": "raw_1c_business_events",
"document_changes": "raw_1c_document_changes",
"companies": "raw_1c_companies",
"reglog": "raw_reglog",
"audit": "raw_audit",
"host": "raw_host_metrics",
}
CORE_TABLES = {
"documents": "documents",
"postings": "postings",
"business_events": "business_events",
"document_changes": "document_change_events",
"companies": "companies",
"reglog": "reglog_events",
"audit": "audit_events",
"host": "host_events",
}
@dataclass
class Config:
clickhouse: dict[str, Any]
landing: dict[str, str]
formats: dict[str, str]
archive_dir: str | None
delete_after_load: bool
min_file_age_seconds: int
def parse_args() -> argparse.Namespace:
p = argparse.ArgumentParser(description="Load file-based 1C exports into ClickHouse")
p.add_argument("--config", required=True, help="Path to YAML config")
p.add_argument("--dataset", choices=["documents", "postings", "business_events", "document_changes", "companies", "reglog", "audit", "host"], help="Load only one dataset")
return p.parse_args()
def load_config(path: str) -> Config:
raw = yaml.safe_load(Path(path).read_text(encoding="utf-8"))
return Config(
clickhouse=raw["clickhouse"],
landing=raw["landing"],
formats=raw.get("formats", {"default": "jsonl"}),
archive_dir=raw.get("archive_dir"),
delete_after_load=bool(raw.get("delete_after_load", False)),
min_file_age_seconds=int(raw.get("min_file_age_seconds", 180)),
)
def ch_client(conf: Config):
import clickhouse_connect
return clickhouse_connect.get_client(
host=conf.clickhouse["host"],
port=conf.clickhouse.get("port", 8123),
username=conf.clickhouse.get("username", "default"),
password=conf.clickhouse.get("password", ""),
database=conf.clickhouse.get("database", "analytics_1c"),
)
def normalize_ts(value: Any) -> datetime:
if isinstance(value, datetime):
return value
if value is None or value == "":
return datetime.utcnow()
return date_parser.parse(str(value))
def iter_rows(path: Path, fmt: str) -> list[dict[str, Any]]:
if fmt == "jsonl":
return [json.loads(line) for line in path.read_text(encoding="utf-8-sig").splitlines() if line.strip()]
if fmt == "json":
payload = json.loads(path.read_text(encoding="utf-8-sig"))
return payload if isinstance(payload, list) else [payload]
if fmt == "csv":
with path.open("r", encoding="utf-8-sig", newline="") as fh:
return list(csv.DictReader(fh))
raise ValueError(f"unsupported format: {fmt}")
def insert_raw(client, dataset: str, source_file: str, rows: list[dict[str, Any]]) -> None:
client.insert(
RAW_TABLES[dataset],
[[source_file, json.dumps(row, ensure_ascii=False)] for row in rows],
column_names=["source_file", "payload"],
)
def map_core_row(dataset: str, source_file: str, row: dict[str, Any]) -> list[Any]:
if dataset == "documents":
return [
normalize_ts(row.get("ts") or row.get("posted_at") or row.get("created_at")),
row.get("infobase", ""),
row.get("organization", ""),
row.get("department", ""),
row.get("doc_type", ""),
row.get("doc_id", ""),
row.get("doc_number", ""),
row.get("author", ""),
row.get("counterparty", ""),
row.get("operation_type", ""),
float(row.get("amount", 0) or 0),
row.get("status", ""),
int(row.get("posted", 0) or 0),
source_file,
]
if dataset == "postings":
return [
normalize_ts(row.get("ts")),
row.get("infobase", ""),
row.get("registrar", ""),
row.get("operation_type", ""),
row.get("account_dt", ""),
row.get("account_ct", ""),
float(row.get("amount", 0) or 0),
source_file,
]
if dataset == "business_events":
return [
normalize_ts(row.get("ts") or row.get("event_time")),
row.get("event_id", ""),
row.get("infobase", ""),
row.get("company_entity_key", ""),
row.get("organization", ""),
row.get("department", ""),
row.get("document_id", row.get("doc_id", "")),
row.get("document_number", row.get("doc_number", "")),
row.get("document_type", row.get("doc_type", "")),
row.get("registrar", ""),
row.get("operation_type", ""),
row.get("event_kind", ""),
row.get("user", row.get("author", "")),
row.get("counterparty", ""),
row.get("counterparty_inn", ""),
row.get("debit_account", row.get("account_dt", "")),
row.get("credit_account", row.get("account_ct", "")),
float(row.get("amount", 0) or 0),
row.get("currency", "RUB"),
int(row.get("line_no", 0) or 0),
row.get("evidence_ref", ""),
source_file,
]
if dataset == "document_changes":
return [
normalize_ts(row.get("ts") or row.get("change_time")),
row.get("change_id", ""),
row.get("infobase", ""),
row.get("company_entity_key", ""),
row.get("organization", ""),
row.get("document_id", row.get("doc_id", "")),
row.get("document_number", row.get("doc_number", "")),
row.get("document_type", row.get("doc_type", "")),
row.get("change_kind", ""),
row.get("field_name", ""),
row.get("user", row.get("author", "")),
row.get("before_value", ""),
row.get("after_value", ""),
row.get("risk_tag", ""),
row.get("evidence_ref", ""),
source_file,
]
if dataset == "companies":
return [
normalize_ts(row.get("ts")),
row.get("infobase", ""),
row.get("company_name", row.get("counterparty", row.get("infobase", ""))),
row.get("organization", ""),
row.get("owner_user", row.get("author", "")),
row.get("base_id", row.get("doc_id", "")),
row.get("base_path", ""),
row.get("status", ""),
int(row.get("db_size_bytes", 0) or 0),
int(row.get("reglog_size_bytes", 0) or 0),
int(row.get("active_locks", 0) or 0),
int(row.get("temp_db_present", 0) or 0),
int(row.get("scheduler_touched", 0) or 0),
float(row.get("activity_score", row.get("amount", 0)) or 0),
source_file,
]
if dataset == "reglog":
return [
normalize_ts(row.get("ts")),
row.get("infobase", ""),
row.get("user", ""),
row.get("host", ""),
row.get("app", ""),
row.get("event_name", ""),
row.get("level", "info"),
int(row.get("duration_ms", 0) or 0),
row.get("message", ""),
source_file,
]
if dataset == "audit":
return [
normalize_ts(row.get("ts")),
row.get("infobase", ""),
row.get("user", ""),
row.get("object_type", ""),
row.get("object_id", ""),
row.get("action", ""),
row.get("before_hash", ""),
row.get("after_hash", ""),
row.get("risk_tag", ""),
source_file,
]
if dataset == "host":
return [
normalize_ts(row.get("ts")),
row.get("host", ""),
float(row.get("cpu_pct", 0) or 0),
float(row.get("ram_pct", 0) or 0),
float(row.get("disk_free_gb", 0) or 0),
float(row.get("disk_latency_ms", 0) or 0),
int(row.get("smb_errors", 0) or 0),
int(row.get("rdp_sessions", 0) or 0),
int(row.get("backup_ok", 0) or 0),
source_file,
]
raise ValueError(dataset)
def core_columns(dataset: str) -> list[str]:
if dataset == "documents":
return ["ts", "infobase", "organization", "department", "doc_type", "doc_id", "doc_number", "author", "counterparty", "operation_type", "amount", "status", "posted", "source_file"]
if dataset == "postings":
return ["ts", "infobase", "registrar", "operation_type", "account_dt", "account_ct", "amount", "source_file"]
if dataset == "business_events":
return ["ts", "event_id", "infobase", "company_entity_key", "organization", "department", "document_id", "document_number", "document_type", "registrar", "operation_type", "event_kind", "user", "counterparty", "counterparty_inn", "debit_account", "credit_account", "amount", "currency", "line_no", "evidence_ref", "source_file"]
if dataset == "document_changes":
return ["ts", "change_id", "infobase", "company_entity_key", "organization", "document_id", "document_number", "document_type", "change_kind", "field_name", "user", "before_value", "after_value", "risk_tag", "evidence_ref", "source_file"]
if dataset == "companies":
return ["ts", "infobase", "company_name", "organization", "owner_user", "base_id", "base_path", "status", "db_size_bytes", "reglog_size_bytes", "active_locks", "temp_db_present", "scheduler_touched", "activity_score", "source_file"]
if dataset == "reglog":
return ["ts", "infobase", "user", "host", "app", "event_name", "level", "duration_ms", "message", "source_file"]
if dataset == "audit":
return ["ts", "infobase", "user", "object_type", "object_id", "action", "before_hash", "after_hash", "risk_tag", "source_file"]
return ["ts", "host", "cpu_pct", "ram_pct", "disk_free_gb", "disk_latency_ms", "smb_errors", "rdp_sessions", "backup_ok", "source_file"]
def archive_or_delete(conf: Config, dataset: str, path: Path) -> None:
if conf.archive_dir:
archive_root = Path(conf.archive_dir) / dataset
archive_root.mkdir(parents=True, exist_ok=True)
shutil.move(str(path), archive_root / path.name)
return
if conf.delete_after_load:
path.unlink(missing_ok=True)
def main() -> int:
args = parse_args()
conf = load_config(args.config)
client = ch_client(conf)
datasets = [args.dataset] if args.dataset else list(RAW_TABLES)
for dataset in datasets:
landing = Path(conf.landing[dataset])
fmt = conf.formats.get(dataset, conf.formats.get("default", "jsonl"))
if not landing.exists():
continue
for path in sorted(p for p in landing.iterdir() if p.is_file()):
age_seconds = max(0, int((datetime.now(UTC) - datetime.fromtimestamp(path.stat().st_mtime, UTC)).total_seconds()))
if age_seconds < conf.min_file_age_seconds:
print(f"skip {dataset}: {path.name} age={age_seconds}s < min_file_age_seconds={conf.min_file_age_seconds}")
continue
rows = iter_rows(path, fmt)
if not rows:
archive_or_delete(conf, dataset, path)
continue
insert_raw(client, dataset, path.name, rows)
client.insert(
CORE_TABLES[dataset],
[map_core_row(dataset, path.name, row) for row in rows],
column_names=core_columns(dataset),
)
archive_or_delete(conf, dataset, path)
print(f"loaded {dataset}: {path.name} rows={len(rows)}")
return 0
if __name__ == "__main__":
raise SystemExit(main())