# DeciKG Demo Notebook (Neo4j原因追跡 × Postgres数値検証)
目的:
- KPIの異常(月次の10月落ち込み)を「グラフで固定」する
- Neo4jで原因候補をランキングし、探索キー(material/process/step)を抽出
- Postgresで数値(単価・使用量・工数)を検証して「証拠」を出す
- すべてのクエリと結果を監査ログに残し、再現手順として残す
Block 1: セットアップ(Imports / パラメータ / .env)(コードセル)
# !pip install pandas sqlalchemy psycopg2-binary neo4j python-dotenv matplotlib
import os, json, hashlib, platform
from datetime import datetime, timezone
from dataclasses import dataclass
from typing import Any, Dict, List, Optional
import pandas as pd
import matplotlib.pyplot as plt
from sqlalchemy import create_engine, text
from neo4j import GraphDatabase
try:
from dotenv import load_dotenv
load_dotenv() # .env を読む(Notebookに接続情報を直書きしない)
except Exception:
pass
# ===== 入力パラメータ(デモで差し替えるのはここだけ)=====
company_key = "demo"
periods = ["2025-08","2025-09","2025-10","2025-11"]
target_kpi = "gross_margin_rate"
# ===== 接続情報(.env から読む想定)=====
POSTGRES_URL="postgresql+psycopg2://odoo:odoo@odoo-db:5432/kpi"
NEO4J_URI="neo4j://neo4j:7687"
NEO4J_USER="neo4j"
NEO4J_PASS="neo4j_pass_123"
PG_SCHEMA="kg_poc"
assert POSTGRES_URL, "POSTGRES_URL が未設定です (.env推奨)"
assert NEO4J_URI and NEO4J_USER and NEO4J_PASS, "NEO4J_URI/NEO4J_USER/NEO4J_PASS が未設定です"
Block 2: 接続(Postgres / Neo4j)(コードセル)
pg_engine = create_engine(POSTGRES_URL, future=True)
neo_driver = GraphDatabase.driver(
NEO4J_URI,
auth=(NEO4J_USER, NEO4J_PASS),
)
print("Connected:", "Postgres OK / Neo4j OK")
Block 3: 監査ログ(Audit Log)基盤(コードセル)
「誰が見ても同じ結果」を担保するために、実行したSQL/Cypherと行数、パラメータ、時刻を JSONL に残します。
AUDIT_PATH = "audit_log.jsonl"
def _now_iso():
return datetime.now(timezone.utc).isoformat()
def _hash_params(obj: Any) -> str:
raw = json.dumps(obj, ensure_ascii=False, sort_keys=True).encode("utf-8")
return hashlib.sha256(raw).hexdigest()
RUN_META = {
"run_id": _hash_params({"company_key": company_key, "periods": periods, "target_kpi": target_kpi, "ts": _now_iso()}),
"ts_utc": _now_iso(),
"python": platform.python_version(),
"platform": platform.platform(),
"params": {"company_key": company_key, "periods": periods, "target_kpi": target_kpi},
}
with open(AUDIT_PATH, "a", encoding="utf-8") as f:
f.write(json.dumps({"type":"run_start", **RUN_META}, ensure_ascii=False) + "\n")
print("audit run_id:", RUN_META["run_id"])
Block 4: 実行ヘルパ(pg_df / neo_df)+ログ記録(コードセル)
def audit_event(event: Dict[str, Any]):
event = {**event, "ts_utc": _now_iso(), "run_id": RUN_META["run_id"]}
with open(AUDIT_PATH, "a", encoding="utf-8") as f:
f.write(json.dumps(event, ensure_ascii=False) + "\n")
def pg_df(sql: str, params: Optional[Dict[str, Any]] = None) -> pd.DataFrame:
params = params or {}
audit_event({"type":"postgres_query", "sql": sql, "params": params})
with pg_engine.connect() as conn:
df = pd.read_sql(text(sql), conn, params=params)
audit_event({"type":"postgres_result", "rows": len(df), "cols": list(df.columns)})
return df
def neo_df(cypher: str, params: Optional[Dict[str, Any]] = None) -> pd.DataFrame:
params = params or {}
audit_event({"type":"neo4j_query", "cypher": cypher, "params": params})
with neo_driver.session() as session:
rows = session.run(cypher, params).data()
df = pd.DataFrame(rows)
audit_event({"type":"neo4j_result", "rows": len(df), "cols": list(df.columns)})
return df
Block 5: KPI推移(Postgres)→ DataFrame(コードセル)
sql_kpi_trend = f"""
SELECT
p.period_key,
to_char(p.start_date,'YYYY-MM-DD') || '〜' || to_char(p.end_date,'YYYY-MM-DD') AS label,
r.value AS kpi_value
FROM {PG_SCHEMA}.kpi_result r
JOIN {PG_SCHEMA}.period p ON p.id = r.period_id
JOIN {PG_SCHEMA}.company c ON c.id = r.company_id
WHERE c.company_key = :company_key
AND r.kpi_key = :kpi_key
AND p.period_key = ANY(:periods)
ORDER BY p.period_key;
"""
df_kpi = pg_df(sql_kpi_trend, {"company_key": company_key, "kpi_key": target_kpi, "periods": periods})
df_kpi
Block 6: KPI推移グラフ(matplotlib)—「BI代替」(コードセル)
plt.figure()
plt.plot(df_kpi["period_key"], df_kpi["kpi_value"], marker="o")
plt.title(f"{company_key} / {target_kpi} trend")
plt.xlabel("period")
plt.ylabel(target_kpi)
plt.grid(True)
plt.show()
# 異常月(最小値)を自動選択(デモだと 2025-10 になる想定)
focus_period = df_kpi.loc[df_kpi["kpi_value"].idxmin(), "period_key"]
print("focus_period =", focus_period)
Block 7: 原因候補ランキング(Neo4j: CONTRIBUTES_TO 上位)(コードセル)
※あなたのNeo4jのラベル/リレーション名が微妙に違う可能性があるので、まずはこの想定Cypherで動かし、0件なら「スキーマ探索Block」を追加して合わせ込みます。
cypher_top_components = """
MATCH (kr:KPIResult {
company_key: $company_key,
kpi_key: $kpi_key,
period_key: $period_key
})
MATCH (kr)-[r:CONTRIBUTES_TO]->(comp:Component)
RETURN
comp.component_key AS component_key,
coalesce(r.contribution_score, r.score, r.weight, 0) AS contribution_score,
coalesce(r.notes, '') AS notes
ORDER BY contribution_score DESC
LIMIT 10
"""
df_top = neo_df(cypher_top_components, {
"company_key": company_key,
"kpi_key": target_kpi,
"period_key": focus_period
})
df_top
Block 8: フォールバック(Neo4jが0件ならPostgres寄与度から取る)(コードセル)
if df_top.empty:
sql_contrib = f"""
SELECT
component_key,
contribution_score,
contribution_type,
notes
FROM {PG_SCHEMA}.kpi_component_contribution kc
JOIN {PG_SCHEMA}.period p ON p.id = kc.period_id
JOIN {PG_SCHEMA}.company c ON c.id = kc.company_id
WHERE c.company_key = :company_key
AND kc.kpi_key = :kpi_key
AND p.period_key = :period_key
ORDER BY kc.contribution_score ASC
LIMIT 10;
"""
df_top = pg_df(sql_contrib, {
"company_key": company_key,
"kpi_key": target_kpi,
"period_key": focus_period
})
df_top.head(3)
Block 9: 上位3コンポーネントを変数で受け渡し(コードセル)
top3 = df_top.head(3).copy()
top_components = top3["component_key"].tolist()
display(top3)
print("top_components =", top_components)
Block 10: Neo4jで探索キー抽出(material_key / process_key / step_key)(コードセル)
ここが「原因追跡(グラフ探索)」です。
コンポーネントごとに、ドライバー(材料/工程/ステップ)を取り出します。
def extract_keys_from_neo4j(component_key: str) -> Dict[str, List[str]]:
# 代表的な関係名を OR で拾う(環境差吸収)
cypher = """
MATCH (comp:Component {component_key:$component_key})
OPTIONAL MATCH (comp)-[:DRIVEN_BY_MATERIAL|:MATERIAL_DRIVER|:USES_MATERIAL]->(m:Material)
OPTIONAL MATCH (comp)-[:DRIVEN_BY_PROCESS|:PROCESS_DRIVER|:USES_PROCESS]->(pr:Process)
OPTIONAL MATCH (pr)-[:HAS_STEP|:PROCESS_STEP|:CONTAINS_STEP]->(st:Step)
RETURN
collect(distinct m.material_key) AS material_keys,
collect(distinct pr.process_key) AS process_keys,
collect(distinct st.step_key) AS step_keys
"""
df = neo_df(cypher, {"component_key": component_key})
if df.empty:
return {"material_keys": [], "process_keys": [], "step_keys": []}
row = df.iloc[0].to_dict()
# None除去
return {
"material_keys": [x for x in row.get("material_keys", []) if x],
"process_keys": [x for x in row.get("process_keys", []) if x],
"step_keys": [x for x in row.get("step_keys", []) if x],
}
keys_by_component = {ck: extract_keys_from_neo4j(ck) for ck in top_components}
keys_by_component
Block 11: Postgresで数値検証(単価・使用量・工数)(コードセル)
あなたが追加した最小Factテーブル(fact_material_price_monthly / fact_material_usage_monthly / fact_labor_hours_monthly)を使います。
def verify_material_price(material_key: str) -> pd.DataFrame:
sql = f"""
SELECT p.period_key, f.material_key, f.avg_unit_price, f.currency, f.notes
FROM {PG_SCHEMA}.fact_material_price_monthly f
JOIN {PG_SCHEMA}.period p ON p.id=f.period_id
JOIN {PG_SCHEMA}.company c ON c.id=f.company_id
WHERE c.company_key=:company_key
AND p.period_key = ANY(:periods)
AND f.material_key=:material_key
ORDER BY p.period_key;
"""
df = pg_df(sql, {"company_key": company_key, "periods": periods, "material_key": material_key})
if not df.empty:
df["mom_delta_price"] = df["avg_unit_price"].diff()
return df
def verify_material_usage(material_key: str) -> pd.DataFrame:
sql = f"""
SELECT p.period_key, f.step_key, f.material_key, f.usage_qty, f.uom, f.notes
FROM {PG_SCHEMA}.fact_material_usage_monthly f
JOIN {PG_SCHEMA}.period p ON p.id=f.period_id
JOIN {PG_SCHEMA}.company c ON c.id=f.company_id
WHERE c.company_key=:company_key
AND p.period_key = ANY(:periods)
AND f.material_key=:material_key
ORDER BY p.period_key;
"""
df = pg_df(sql, {"company_key": company_key, "periods": periods, "material_key": material_key})
if not df.empty:
df["mom_delta_usage"] = df["usage_qty"].diff()
return df
def verify_labor_hours(process_key: Optional[str] = None, step_key: Optional[str] = None) -> pd.DataFrame:
sql = f"""
SELECT p.period_key, f.process_key, f.step_key, f.skill_key, f.labor_hours, f.notes
FROM {PG_SCHEMA}.fact_labor_hours_monthly f
JOIN {PG_SCHEMA}.period p ON p.id=f.period_id
JOIN {PG_SCHEMA}.company c ON c.id=f.company_id
WHERE c.company_key=:company_key
AND p.period_key = ANY(:periods)
AND (:process_key IS NULL OR f.process_key=:process_key)
AND (:step_key IS NULL OR f.step_key=:step_key)
ORDER BY p.period_key;
"""
df = pg_df(sql, {"company_key": company_key, "periods": periods, "process_key": process_key, "step_key": step_key})
if not df.empty:
df["mom_delta_hours"] = df["labor_hours"].diff()
return df
# 代表例:抽出キーの先頭だけ検証して表示(デモはここを見せ場に)
material_candidates = []
process_candidates = []
for ck, keys in keys_by_component.items():
material_candidates += keys["material_keys"]
process_candidates += keys["process_keys"]
material_candidates = sorted(set([m for m in material_candidates if m]))
process_candidates = sorted(set([p for p in process_candidates if p]))
print("material_candidates:", material_candidates)
print("process_candidates :", process_candidates)
dfs_verify = {}
if material_candidates:
mk = material_candidates[0]
dfs_verify["material_price"] = verify_material_price(mk)
dfs_verify["material_usage"] = verify_material_usage(mk)
if process_candidates:
pk = process_candidates[0]
dfs_verify["labor_hours"] = verify_labor_hours(process_key=pk)
dfs_verify
Block 12: Evidence(証拠)を添える(Postgres)(コードセル)
sql_evidence = f"""
SELECT
e.evidence_key,
e.summary,
e.source_system,
e.source_ref
FROM {PG_SCHEMA}.kpi_result r
JOIN {PG_SCHEMA}.period p ON p.id=r.period_id
JOIN {PG_SCHEMA}.company c ON c.id=r.company_id
LEFT JOIN {PG_SCHEMA}.kpi_evidence_map kem ON kem.kpi_result_id = r.id
LEFT JOIN {PG_SCHEMA}.evidence e ON (
(pg_temp.col_exists('{PG_SCHEMA}','kpi_evidence_map','evidence_id') AND e.id = kem.evidence_id)
OR (pg_temp.col_exists('{PG_SCHEMA}','kpi_evidence_map','evidence_key') AND e.evidence_key = kem.evidence_key)
)
WHERE c.company_key=:company_key
AND r.kpi_key=:kpi_key
AND p.period_key=:period_key;
"""
# ※上のSQLは pg_temp 関数をSQL側で呼ぶ形なので、環境によっては嫌がる場合があります。
# その場合は「evidence_id列がある/ない」の分岐をPython側で切り替える版にします。
# まずは "evidence_id/evidence_key どっちの列があるか" をPython側で確認
col_check = pg_df(f"""
SELECT
EXISTS(SELECT 1 FROM information_schema.columns WHERE table_schema='{PG_SCHEMA}' AND table_name='kpi_evidence_map' AND column_name='evidence_id') AS has_evidence_id,
EXISTS(SELECT 1 FROM information_schema.columns WHERE table_schema='{PG_SCHEMA}' AND table_name='kpi_evidence_map' AND column_name='evidence_key') AS has_evidence_key
""")
has_id = bool(col_check.loc[0,"has_evidence_id"])
has_key = bool(col_check.loc[0,"has_evidence_key"])
if has_id:
sql_evidence_py = f"""
SELECT e.evidence_key, e.summary, e.source_system, e.source_ref
FROM {PG_SCHEMA}.kpi_result r
JOIN {PG_SCHEMA}.period p ON p.id=r.period_id
JOIN {PG_SCHEMA}.company c ON c.id=r.company_id
JOIN {PG_SCHEMA}.kpi_evidence_map kem ON kem.kpi_result_id=r.id
JOIN {PG_SCHEMA}.evidence e ON e.id=kem.evidence_id
WHERE c.company_key=:company_key
AND r.kpi_key=:kpi_key
AND p.period_key=:period_key;
"""
elif has_key:
sql_evidence_py = f"""
SELECT e.evidence_key, e.summary, e.source_system, e.source_ref
FROM {PG_SCHEMA}.kpi_result r
JOIN {PG_SCHEMA}.period p ON p.id=r.period_id
JOIN {PG_SCHEMA}.company c ON c.id=r.company_id
JOIN {PG_SCHEMA}.kpi_evidence_map kem ON kem.kpi_result_id=r.id
JOIN {PG_SCHEMA}.evidence e ON e.evidence_key=kem.evidence_key
WHERE c.company_key=:company_key
AND r.kpi_key=:kpi_key
AND p.period_key=:period_key;
"""
else:
sql_evidence_py = None
df_evid = pd.DataFrame()
if sql_evidence_py:
df_evid = pg_df(sql_evidence_py, {"company_key": company_key, "kpi_key": target_kpi, "period_key": focus_period})
df_evid
Block 13: 結論(Markdownセル)
Notebookの最後に「一枚で結論」を置きます。数字は上の DataFrame から拾って文章化します(最初は固定文でもOK)。
## 結論(Explain / 監査用)
- 10月(focus_period)に粗利益率が最も低下したことをKPI推移で確認
- 原因候補ランキング(上位3コンポーネント)を取得
- 主要因の探索キー(material_key/process_key/step_key)を抽出
- Postgresで「単価・使用量・工数」の推移を検証し、前月差が異常であることを確認
- Evidence(例: EVID-STEEL-PRICE-2025-10)を紐付けて説明可能性を補強
例:
「10月GM低下は、鋼板単価↑と溶接工数↑が主要因。証拠は EVID-STEEL-PRICE-2025-10」
Block 14: 重要セルを関数化(run_demo)(コードセル)
「テンプレ化して顧客ごとに差し替え」を前提に、最後に run_demo() を用意します。
def run_demo(company_key_: str, periods_: List[str], target_kpi_: str) -> Dict[str, Any]:
# (Notebookの流れを最低限まとめた版)
out: Dict[str, Any] = {}
out["params"] = {"company_key": company_key_, "periods": periods_, "target_kpi": target_kpi_}
df_kpi_ = pg_df(sql_kpi_trend, {"company_key": company_key_, "kpi_key": target_kpi_, "periods": periods_})
out["kpi_trend"] = df_kpi_
focus_ = df_kpi_.loc[df_kpi_["kpi_value"].idxmin(), "period_key"]
out["focus_period"] = focus_
# contribution: neo->fallback pg
df_top_ = neo_df(cypher_top_components, {"company_key": company_key_, "kpi_key": target_kpi_, "period_key": focus_})
if df_top_.empty:
df_top_ = pg_df(f"""
SELECT component_key, contribution_score, contribution_type, notes
FROM {PG_SCHEMA}.kpi_component_contribution kc
JOIN {PG_SCHEMA}.period p ON p.id = kc.period_id
JOIN {PG_SCHEMA}.company c ON c.id = kc.company_id
WHERE c.company_key = :company_key
AND kc.kpi_key = :kpi_key
AND p.period_key = :period_key
ORDER BY kc.contribution_score ASC
LIMIT 10;
""", {"company_key": company_key_, "kpi_key": target_kpi_, "period_key": focus_})
out["top_components"] = df_top_.head(3)
top_components_ = out["top_components"]["component_key"].tolist()
keys_ = {ck: extract_keys_from_neo4j(ck) for ck in top_components_}
out["keys_by_component"] = keys_
# 代表キーだけ数値検証
mats = sorted(set(sum([v["material_keys"] for v in keys_.values()], [])))
procs = sorted(set(sum([v["process_keys"] for v in keys_.values()], [])))
if mats:
out["verify_material_price"] = verify_material_price(mats[0])
out["verify_material_usage"] = verify_material_usage(mats[0])
if procs:
out["verify_labor_hours"] = verify_labor_hours(process_key=procs[0])
return out
result = run_demo(company_key, periods, target_kpi)
result.keys()
コメントを残す