feat(1c): scaffold business event extraction layer

This commit is contained in:
igor04091968
2026-05-22 21:47:52 +03:00
parent 6b246aaa14
commit 61770c1658
8 changed files with 420 additions and 2 deletions
+18 -1
View File
@@ -26,6 +26,8 @@ File 1C + reglog + host telemetry
├─ raw_*
├─ documents
├─ postings
├─ business_events
├─ document_change_events
├─ reglog_events
├─ audit_events
├─ host_events
@@ -45,6 +47,8 @@ File 1C + reglog + host telemetry
- `clickhouse/init/*.sql` — схема БД.
- `etl/load_1c_exports.py` — loader CSV/JSON выгрузок в raw/core таблицы.
- `etl/config.example.yml` — пример ETL-конфига.
- `docs/1C_BUSINESS_EVENT_LAYER_RU.md` — production contract следующего шага:
canonical business-event слой для документов/проводок/изменений.
- `detections/rules.yml` — каталог правил detections.
- `detections/insert_detections.sql` — SQL-шаблоны rule-based detections.
- `grafana/dashboard-catalog.md` — целевая структура дашбордов.
@@ -95,7 +99,7 @@ docker compose up -d
3. Инициализировать landing-каталоги и ETL config:
```bash
mkdir -p landing/{documents,postings,reglog,audit,host}
mkdir -p landing/{documents,postings,business_events,document_changes,companies,reglog,audit,host}
cp etl/config.example.yml etl/config.yml
```
@@ -156,6 +160,8 @@ Read-only API для руководителя:
- выгрузки 1С по документам;
- выгрузки движений/проводок;
- выгрузки canonical business events;
- выгрузки изменений документов и реквизитов;
- read-only `companies` snapshot по файловым базам;
- журнал регистрации 1С;
- audit/export критичных изменений;
@@ -184,3 +190,14 @@ Read-only API для руководителя:
- case/timeline слой считается вне 1С.
- это прогноз активности компании/базы, а не финансовых проводок;
- если `counterparty` в live-выгрузках пустой, company-forecast слой останется корректно пустым.
## Следующий production шаг
В репо уже заложен scaffold под следующий слой:
- `business_events` — единый event stream бухгалтерских событий;
- `document_change_events` — изменения документов/реквизитов;
- ETL уже умеет принимать эти datasets в `landing/business_events` и
`landing/document_changes`;
- дальше нужен только read-only extractor, который будет наполнять их из 1С
или внешних безопасных выгрузок без записи в базу.
@@ -16,6 +16,24 @@ CREATE TABLE IF NOT EXISTS analytics_1c.raw_1c_postings
ENGINE = MergeTree
ORDER BY (ingested_at, source_file);
CREATE TABLE IF NOT EXISTS analytics_1c.raw_1c_business_events
(
ingested_at DateTime DEFAULT now(),
source_file String,
payload String
)
ENGINE = MergeTree
ORDER BY (ingested_at, source_file);
CREATE TABLE IF NOT EXISTS analytics_1c.raw_1c_document_changes
(
ingested_at DateTime DEFAULT now(),
source_file String,
payload String
)
ENGINE = MergeTree
ORDER BY (ingested_at, source_file);
CREATE TABLE IF NOT EXISTS analytics_1c.raw_1c_companies
(
ingested_at DateTime DEFAULT now(),
@@ -32,6 +32,56 @@ CREATE TABLE IF NOT EXISTS analytics_1c.postings
ENGINE = MergeTree
ORDER BY (infobase, ts, registrar);
CREATE TABLE IF NOT EXISTS analytics_1c.business_events
(
ts DateTime,
event_id String,
infobase LowCardinality(String),
company_entity_key String,
organization String,
department String,
document_id String,
document_number String,
document_type LowCardinality(String),
registrar String,
operation_type String,
event_kind LowCardinality(String),
user String,
counterparty String,
counterparty_inn String,
debit_account String,
credit_account String,
amount Decimal(18, 2),
currency LowCardinality(String),
line_no UInt32,
evidence_ref String,
source_file String
)
ENGINE = MergeTree
ORDER BY (infobase, ts, document_id, line_no, event_id);
CREATE TABLE IF NOT EXISTS analytics_1c.document_change_events
(
ts DateTime,
change_id String,
infobase LowCardinality(String),
company_entity_key String,
organization String,
document_id String,
document_number String,
document_type LowCardinality(String),
change_kind LowCardinality(String),
field_name String,
user String,
before_value String,
after_value String,
risk_tag String,
evidence_ref String,
source_file String
)
ENGINE = MergeTree
ORDER BY (infobase, ts, document_id, change_id);
CREATE TABLE IF NOT EXISTS analytics_1c.companies
(
ts DateTime,
+4
View File
@@ -8,6 +8,8 @@ clickhouse:
landing:
documents: ./landing/documents
postings: ./landing/postings
business_events: ./landing/business_events
document_changes: ./landing/document_changes
companies: ./landing/companies
reglog: ./landing/reglog
audit: ./landing/audit
@@ -17,6 +19,8 @@ formats:
default: jsonl
documents: jsonl
postings: jsonl
business_events: jsonl
document_changes: jsonl
companies: jsonl
reglog: jsonl
audit: jsonl
+53 -1
View File
@@ -17,6 +17,8 @@ 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",
@@ -26,6 +28,8 @@ RAW_TABLES = {
CORE_TABLES = {
"documents": "documents",
"postings": "postings",
"business_events": "business_events",
"document_changes": "document_change_events",
"companies": "companies",
"reglog": "reglog_events",
"audit": "audit_events",
@@ -46,7 +50,7 @@ class Config:
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", "companies", "reglog", "audit", "host"], help="Load only one dataset")
p.add_argument("--dataset", choices=["documents", "postings", "business_events", "document_changes", "companies", "reglog", "audit", "host"], help="Load only one dataset")
return p.parse_args()
@@ -129,6 +133,50 @@ def map_core_row(dataset: str, source_file: str, row: dict[str, Any]) -> list[An
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")),
@@ -194,6 +242,10 @@ def core_columns(dataset: str) -> list[str]:
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":
+81
View File
@@ -0,0 +1,81 @@
#!/usr/bin/env python3
from __future__ import annotations
import sys
import unittest
from pathlib import Path
from types import SimpleNamespace
sys.path.insert(0, str(Path(__file__).resolve().parent))
sys.modules.setdefault("clickhouse_connect", SimpleNamespace(get_client=lambda **_: None))
import load_1c_exports as loader
class Load1CExportsTests(unittest.TestCase):
def test_business_events_mapping(self) -> None:
row = {
"event_time": "2026-05-22T12:00:00Z",
"event_id": "evt-1",
"infobase": "ИБ-1",
"company_entity_key": "baseid:abc",
"organization": "ООО Тест",
"department": "Продажи",
"doc_id": "doc-1",
"doc_number": "0001",
"doc_type": "Реализация",
"registrar": "DOC-1",
"operation_type": "sale",
"event_kind": "posting",
"author": "user1",
"counterparty": "Контрагент",
"counterparty_inn": "7700000000",
"account_dt": "62.01",
"account_ct": "90.01",
"amount": "125000.50",
"line_no": "2",
"evidence_ref": "reglog:1",
}
mapped = loader.map_core_row("business_events", "events.jsonl", row)
self.assertEqual(mapped[1], "evt-1")
self.assertEqual(mapped[3], "baseid:abc")
self.assertEqual(mapped[7], "0001")
self.assertEqual(mapped[15], "62.01")
self.assertEqual(mapped[17], 125000.5)
self.assertEqual(mapped[19], 2)
self.assertEqual(mapped[21], "events.jsonl")
def test_document_changes_mapping(self) -> None:
row = {
"change_time": "2026-05-22T12:00:00Z",
"change_id": "chg-1",
"infobase": "ИБ-1",
"company_entity_key": "basepath:c:/1c/base",
"organization": "ООО Тест",
"document_id": "doc-2",
"document_number": "0002",
"document_type": "Поступление",
"change_kind": "requisites_change",
"field_name": "Контрагент",
"user": "user2",
"before_value": "Старый",
"after_value": "Новый",
"risk_tag": "counterparty_change",
"evidence_ref": "audit:2",
}
mapped = loader.map_core_row("document_changes", "changes.jsonl", row)
self.assertEqual(mapped[1], "chg-1")
self.assertEqual(mapped[3], "basepath:c:/1c/base")
self.assertEqual(mapped[8], "requisites_change")
self.assertEqual(mapped[12], "Новый")
self.assertEqual(mapped[15], "changes.jsonl")
def test_new_dataset_columns_shape(self) -> None:
self.assertEqual(loader.RAW_TABLES["business_events"], "raw_1c_business_events")
self.assertEqual(loader.CORE_TABLES["document_changes"], "document_change_events")
self.assertEqual(len(loader.core_columns("business_events")), 22)
self.assertEqual(len(loader.core_columns("document_changes")), 16)
if __name__ == "__main__":
unittest.main()
+4
View File
@@ -6,6 +6,8 @@ ROOT="${AW_1C_ROOT:-/opt/activitywatch/clickhouse-1c}"
mkdir -p \
"${ROOT}/landing/documents" \
"${ROOT}/landing/postings" \
"${ROOT}/landing/business_events" \
"${ROOT}/landing/document_changes" \
"${ROOT}/landing/companies" \
"${ROOT}/landing/registry" \
"${ROOT}/landing/reglog" \
@@ -13,6 +15,8 @@ mkdir -p \
"${ROOT}/landing/host" \
"${ROOT}/archive/documents" \
"${ROOT}/archive/postings" \
"${ROOT}/archive/business_events" \
"${ROOT}/archive/document_changes" \
"${ROOT}/archive/companies" \
"${ROOT}/archive/registry" \
"${ROOT}/archive/reglog" \
+192
View File
@@ -0,0 +1,192 @@
# 1C Business Event Layer
Это следующий production-шаг поверх уже работающего контура
`file 1C -> ClickHouse -> Grafana -> AI Investigator`.
Цель не в том, чтобы заменить текущий `company intelligence`, а в том, чтобы
добавить **read-only business-event слой**, пригодный для финансовых
расследований, explainability и rule-based detections по бухгалтерскому смыслу.
## Почему этот слой нужен
Текущий контур уже решает:
- operational telemetry;
- timeline по reglog/audit;
- detections/cases;
- company portfolio intelligence;
- manager/recovery briefs.
Но он ещё не даёт полноценного ответа на вопросы уровня:
- какие проводки дали вклад в аномалию;
- кто и когда перепровёл документ;
- какие изменения реквизитов повлияли на результат;
- почему вырос НДС, возвраты или нетипичные движения.
Для этого нужен отдельный канонический слой бизнес-событий.
## Слои данных
### 1. `documents`
Карточки документов и базовые агрегаты по документам.
### 2. `postings`
Лёгкий слой проводок/движений. Уже есть в контуре, но он недостаточен как
канонический event stream.
### 3. `business_events`
Новый целевой канонический слой. Один ряд = одно бизнес-событие, пригодное для:
- timeline;
- detections;
- explainability;
- AI investigations;
- correlation с reglog/audit/cases.
Текущая schema scaffold:
- `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`
### 4. `document_change_events`
Новый слой изменений документов и реквизитов.
Нужен для:
- расследования перепроведений;
- контроля изменений реквизитов;
- reconstruction narrative;
- объяснения, какие именно изменения дали бизнес-эффект.
Текущая schema scaffold:
- `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`
## Read-only extraction path
Правильный extraction path такой:
1. Внешний extractor читает только безопасные read-only источники.
2. Формирует `jsonl/csv` в landing-каталоги.
3. `etl/load_1c_exports.py` грузит данные в raw/core ClickHouse tables.
4. Detection/AI слой работает только с ClickHouse.
То есть LLM и manager pages не ходят в 1С напрямую.
## Безопасные источники для extractor
Подходящие:
- регламентированные выгрузки документов/движений;
- журнал регистрации 1С;
- внешние реестры и справочники;
- audit/export критичных изменений;
- отдельные read-only файлы, формируемые рядом с 1С.
Не подходящие по умолчанию:
- запись в файловую базу;
- опасные `COM`/`Configurator` сценарии;
- любой write-back path в production 1С.
## Почему здесь нужен `company_entity_key`
Этот слой не должен зависеть от того, как компания названа в текущий момент в
1С. Поэтому новые event tables уже сразу завязаны на `company_entity_key`.
Приоритет идентификации:
1. `baseid:<base_id>`
2. `basepath:<normalized path>`
3. только fallback на human name
Это делает timeline и расследования устойчивыми к rename.
## Что extractor должен уметь первым
Минимальный read-only extractor v1 должен уметь:
- выгружать документы;
- выгружать проводки;
- выгружать `business_events`;
- выгружать `document_change_events`;
- стабильно наполнять `company_entity_key`;
- писать `evidence_ref`, чтобы расследование не было бездоказательным.
## Первые детекты на этом слое
После появления business-event выгрузок стоит вводить:
- ночные проводки;
- дробление платежей;
- повторные перепроведения;
- изменение реквизитов перед/после движения;
- возвраты после закрытия периода;
- нехарактерные движения по пользователю;
- циклические движения по контрагенту/счётам.
## Что уже сделано в репо
Уже добавлены:
- raw tables:
- `raw_1c_business_events`
- `raw_1c_document_changes`
- core tables:
- `business_events`
- `document_change_events`
- ETL support:
- `landing/business_events`
- `landing/document_changes`
- dataset mapping в `etl/load_1c_exports.py`
Это именно scaffold, а не обещание, что live extractor уже существует.
## Что делать дальше
1. Реализовать read-only extractor в отдельном модуле.
2. Стабильно наполнять `company_entity_key`.
3. Добавить первые SQL detections на `business_events`.
4. Обогащать `entity_timeline` уже не только telemetry/audit, но и
business-event evidence.
5. Включить AI Investigator narrative по реальным проводкам и изменениям.