构建分析应用
适用版本CloudEnterprise
Tier0 使用 Marimo Notebook 通过 Python 进行高级数据分析。本指南以 Bowtie 应用作为示例。
flowchart LR
collect["数据连接"] -->|"采集数据"| uns[("UNS<br/>模型")]
uns -->|"原始数据"| notebook["Notebook"]
notebook -->|"分析结果"| uns
uns -->|"分析结果"| builder["App Builder"]
classDef t0accent fill:#EAF8C8,stroke:#A6CF38,stroke-width:1px,color:#171717
classDef t0soft fill:#F7FAF2,stroke:#D8E6B8,stroke-width:1px,color:#2A2A2A
classDef t0agent fill:#EEF4FF,stroke:#B7C7E8,stroke-width:1px,color:#1F2937
class uns t0accent
class collect,notebook t0soft
class builder t0agent
炼化腐蚀是一个受多种运行因素影响的复杂过程。本示例展示一个实时腐蚀风险评估:结合过程数据估算腐蚀发生概率,并支持主动维护。
获取原始数据
Section titled “获取原始数据”-
在 Tier0 中进入 UNS,导入数据模型。
显示 UNS 模型 JSON
Terminal window {"notes": "type:PATH|TOPIC,topicType:STATE|ACTION|METRIC,schema.type:INTEGER|STRING|FLOAT|DOUBLE|BOOLEAN|LONG|DATETIME","namespace": [{"name": "Refinery","type": "PATH","children": [{"name": "CDU_Plant","type": "PATH","children": [{"name": "State","type": "PATH","topicType": "STATE","children": [{"name": "Corrosion_Risk","type": "TOPIC","topicType": "STATE","description": "Current corrosion risk state inferred by Bayesian Network","enableHistory": "TRUE","schema": [{"name": "risk_state","type": "STRING"},{"name": "previous_state","type": "STRING"},{"name": "model_confidence","type": "FLOAT"},{"name": "timestamp","type": "DATETIME"}]}]},{"name": "Metric","type": "PATH","topicType": "METRIC","children": [{"name": "Corrosion_Monitoring","type": "TOPIC","topicType": "METRIC","description": "Real-time process measurements for CDU atmospheric overhead corrosion monitoring","enableHistory": "TRUE","mockData": "FALSE","schema": [{"name": "d103_chloride","type": "FLOAT"},{"name": "d103_ph","type": "FLOAT"},{"name": "wash_water_flow","type": "FLOAT"},{"name": "desalter_salt_ptb","type": "FLOAT"},{"name": "desalter_bsw","type": "FLOAT"},{"name": "wash_water_rate","type": "FLOAT"},{"name": "rrd_ph","type": "FLOAT"},{"name": "rrd_chloride","type": "FLOAT"},{"name": "rrd_total_iron","type": "FLOAT"},{"name": "sap_evidence","type": "STRING"},{"name": "injection_quill_fail","type": "BOOLEAN"},{"name": "timestamp","type": "DATETIME"}]},{"name": "Corrosion_Analysis","type": "TOPIC","topicType": "METRIC","description": "Posterior probabilities generated by the Bayesian Network","enableHistory": "TRUE","mockData": "FALSE","schema": [{"name": "p_normal","type": "DOUBLE"},{"name": "p_developing","type": "DOUBLE"},{"name": "p_confirmed","type": "DOUBLE"},{"name": "total_corrosion_risk","type": "DOUBLE"},{"name": "lopc_probability","type": "DOUBLE"},{"name": "shutdown_probability","type": "DOUBLE"},{"name": "escalation_probability","type": "DOUBLE"},{"name": "timestamp","type": "DATETIME"}]}]}]}]}]} -
进入 数据流程,创建 数据连接,连接原始数据并发布到 UNS。
显示数据连接 JSON
Terminal window [{"id": "6173469ed5064ccf","type": "tab","label": "corrosion","disabled": false,"info": ""},{"id": "ad20b016baf7247f","type": "OpcUa-Client","z": "6173469ed5064ccf","endpoint": "66bd118c33eb12f8","action": "subscribe","deadbandtype": "a","deadbandvalue": 1,"time": 10,"timeUnit": "s","certificate": "n","localfile": "","localkeyfile": "","securitymode": "None","securitypolicy": "None","useTransport": false,"maxChunkCount": 1,"maxMessageSize": 8192,"receiveBufferSize": 8192,"sendBufferSize": 8192,"setstatusandtime": false,"keepsessionalive": false,"name": "","x": 280,"y": 440,"wires": [["ac9195bb4863dbdd","5a3d61b2ae40dcb6"],[],[]]},{"id": "b6746ed066ba0e1a","type": "inject","z": "6173469ed5064ccf","name": "","props": [{"p": "payload"},{"p": "topic","vt": "str"}],"repeat": "","crontab": "","once": false,"onceDelay": 0.1,"topic": "multiple","payload": "[ { \"name\": \"D103_CHLORIDE\", \"nodeId\": \"ns=2;i=3\", \"datatype\": \"Float\" }, { \"name\": \"D103_PH\", \"nodeId\": \"ns=2;i=4\", \"datatype\": \"Float\" }, { \"name\": \"WASH_WATER_FLOW\", \"nodeId\": \"ns=2;i=5\", \"datatype\": \"Float\" }, { \"name\": \"DESALTER_SALT_PTB\", \"nodeId\": \"ns=2;i=6\", \"datatype\": \"Float\" }, { \"name\": \"DESALTER_BSW\", \"nodeId\": \"ns=2;i=7\", \"datatype\": \"Float\" }, { \"name\": \"WASH_WATER_RATE\", \"nodeId\": \"ns=2;i=8\", \"datatype\": \"Float\" }, { \"name\": \"RRD_PH\", \"nodeId\": \"ns=2;i=9\", \"datatype\": \"Float\" }, { \"name\": \"RRD_CHLORIDE\", \"nodeId\": \"ns=2;i=10\", \"datatype\": \"Float\" }, { \"name\": \"RRD_TOTAL_IRON\", \"nodeId\": \"ns=2;i=11\", \"datatype\": \"Float\" }, { \"name\": \"SAP_EVIDENCE\", \"nodeId\": \"ns=2;i=12\", \"datatype\": \"UInt16\" }, { \"name\": \"INJECTION_QUILL_FAIL\", \"nodeId\": \"ns=2;i=13\", \"datatype\": \"Boolean\" } ]","payloadType": "json","x": 90,"y": 440,"wires": [["ad20b016baf7247f"]]},{"id": "ac9195bb4863dbdd","type": "debug","z": "6173469ed5064ccf","name": "debug 2","active": true,"tosidebar": true,"console": false,"tostatus": false,"complete": "false","statusVal": "","statusType": "auto","x": 490,"y": 480,"wires": []},{"id": "5a3d61b2ae40dcb6","type": "function","z": "6173469ed5064ccf","name": "function 1","func": "const TOPIC = \"Refinery/CDU_Plant/Metric/Corrosion_Monitoring\";\n\nconst nodeMap = {\n \"ns=2;i=3\": { field: \"d103_chloride\", type: \"float\" },\n \"ns=2;i=4\": { field: \"d103_ph\", type: \"float\" },\n \"ns=2;i=5\": { field: \"wash_water_flow\", type: \"float\" },\n \"ns=2;i=6\": { field: \"desalter_salt_ptb\", type: \"float\" },\n \"ns=2;i=7\": { field: \"desalter_bsw\", type: \"float\" },\n \"ns=2;i=8\": { field: \"wash_water_rate\", type: \"float\" },\n \"ns=2;i=9\": { field: \"rrd_ph\", type: \"float\" },\n \"ns=2;i=10\": { field: \"rrd_chloride\", type: \"float\" },\n \"ns=2;i=11\": { field: \"rrd_total_iron\", type: \"float\" },\n \"ns=2;i=12\": { field: \"sap_evidence\", type: \"sap\" },\n \"ns=2;i=13\": { field: \"injection_quill_fail\", type: \"boolean\" }\n};\n\nconst sapMap = {\n 0: \"none\",\n 1: \"work_order\",\n 2: \"inspection_finding\"\n};\n\nfunction normalizeNodeId(v) {\n if (!v) return \"\";\n if (typeof v === \"string\") return v;\n if (v.toString) return v.toString();\n return \"\";\n}\n\nfunction extractValue(payload) {\n // OPC UA DataValue: payload.value.value\n if (payload && payload.value && payload.value.value !== undefined) {\n return payload.value.value;\n }\n\n // Some OPC UA nodes output Variant directly.\n if (payload && payload.value !== undefined) {\n return payload.value;\n }\n\n return payload;\n}\n\nfunction extractNodeId(msg) {\n return normalizeNodeId(\n msg.nodeId ||\n msg.topic ||\n msg.payload?.nodeId ||\n msg.payload?.nodeid\n );\n}\n\nconst nodeId = extractNodeId(msg);\nconst spec = nodeMap[nodeId];\n\nif (!spec) {\n node.warn(`Unknown OPC UA nodeId: ${nodeId}`);\n return null;\n}\n\nlet value = extractValue(msg.payload);\n\nif (spec.type === \"float\") {\n value = Number(value);\n} else if (spec.type === \"boolean\") {\n value = Boolean(value);\n} else if (spec.type === \"sap\") {\n value = sapMap[Number(value)] || \"unknown\";\n}\n\nlet latest = context.get(\"latest\") || {};\nlatest[spec.field] = value;\nlatest.timestamp = new Date().toISOString();\ncontext.set(\"latest\", latest);\n\n// Publish only after all 11 fields have been received at least once.\nconst requiredFields = Object.values(nodeMap).map(x => x.field);\nconst ready = requiredFields.every(field => latest[field] !== undefined);\n\nif (!ready) {\n node.status({\n fill: \"yellow\",\n shape: \"ring\",\n text: `waiting ${Object.keys(latest).length - 1}/11`\n });\n return null;\n}\n\nmsg.topic = TOPIC;\nmsg.payload = {\n d103_chloride: latest.d103_chloride,\n d103_ph: latest.d103_ph,\n wash_water_flow: latest.wash_water_flow,\n desalter_salt_ptb: latest.desalter_salt_ptb,\n desalter_bsw: latest.desalter_bsw,\n wash_water_rate: latest.wash_water_rate,\n rrd_ph: latest.rrd_ph,\n rrd_chloride: latest.rrd_chloride,\n rrd_total_iron: latest.rrd_total_iron,\n sap_evidence: latest.sap_evidence,\n injection_quill_fail: latest.injection_quill_fail,\n timestamp: latest.timestamp\n};\n\nnode.status({\n fill: \"green\",\n shape: \"dot\",\n text: \"published corrosion monitoring\"\n});\n\nreturn msg;","outputs": 1,"timeout": 0,"noerr": 0,"initialize": "","finalize": "","libs": [],"x": 370,"y": 580,"wires": [["aa15be33a84fb812","4fdd68dda3ee5c9e"]]},{"id": "aa15be33a84fb812","type": "mqtt out","z": "6173469ed5064ccf","name": "","topic": "","qos": "","retain": "","respTopic": "","contentType": "","userProps": "","correl": "","expiry": "","broker": "broker-djx75cwgiywq","x": 550,"y": 580,"wires": []},{"id": "4fdd68dda3ee5c9e","type": "debug","z": "6173469ed5064ccf","name": "debug 3","active": true,"tosidebar": true,"console": false,"tostatus": false,"complete": "false","statusVal": "","statusType": "auto","x": 450,"y": 680,"wires": []},{"id": "66bd118c33eb12f8","type": "OpcUa-Endpoint","endpoint": "opc.tcp://172.31.151.237:4850","secpol": "None","secmode": "None","none": true,"login": false,"usercert": false,"usercertificate": "","userprivatekey": ""},{"id": "broker-djx75cwgiywq","type": "mqtt-broker","z": "6173469ed5064ccf","name": "corrosion","broker": "emqx","port": "1883","clientid": "357356352615696","usetls": false,"protocolVersion": "4","keepalive": "60","cleansession": true,"birthTopic": "","birthQos": "0","birthPayload": "","closeTopic": "","closeQos": "0","closePayload": "","willTopic": "","willQos": "0","willPayload": ""},{"id": "c5201f1ad482b182","type": "global-config","env": [],"modules": {"node-red-contrib-opcua": "0.2.339"}}]
在 Notebook 中构建分析应用
Section titled “在 Notebook 中构建分析应用”-
在 Tier0 中进入 Notebook,创建一个新的 notebook。
-
打开该 notebook,添加以下 cell 来分析 UNS 中的数据。
-
单元 1 - 导入所需依赖,并配置 MQTT broker、topics 和运行时设置。
显示代码
Terminal window @app.celldef _():"""Imports and runtime configuration."""import jsonimport mathimport osimport threadingimport timefrom datetime import datetime, timezoneimport marimo as moimport pandas as pdimport paho.mqtt.client as mqttMQTT_BROKER_HOST = os.getenv("TIER0_MQTT_HOST", "enterprise-new-dev.tier0.dev")MQTT_BROKER_PORT = int(os.getenv("TIER0_MQTT_PORT", "1883"))MQTT_USERNAME = os.getenv("TIER0_MQTT_USERNAME", "357890055321856")MQTT_PASSWORD = os.getenv("TIER0_MQTT_PASSWORD", "caa6b34242002a3dff4948dd73ce8a")MQTT_CLIENT_ID = os.getenv("TIER0_MQTT_CLIENT_ID","357890055321856_bn",)SOURCE_TOPIC = "Refinery/CDU_Plant/Metric/Corrosion_Monitoring"ANALYSIS_TOPIC = "Refinery/CDU_Plant/Metric/Corrosion_Analysis"RISK_TOPIC = "Refinery/CDU_Plant/State/Corrosion_Risk"SOURCE_QOS = int(os.getenv("TIER0_SOURCE_QOS", "1"))OUTPUT_QOS = int(os.getenv("TIER0_OUTPUT_QOS", "1"))AUTOSTART = os.getenv("TIER0_BN_AUTOSTART", "true").lower() in {"1", "true", "yes", "on"}ROOT_CONFIDENCE = float(os.getenv("TIER0_BN_ROOT_CONFIDENCE", "0.85"))def utc_now_iso() -> str:return (datetime.now(timezone.utc).isoformat(timespec="milliseconds").replace("+00:00", "Z"))mo.md("""# CDU Corrosion Bayesian AnalysisThis notebook subscribes to a complete corrosion-monitoring snapshot,performs one Bayesian inference for each new source message, and publishesthe analysis results back to the UNS.""")return (ANALYSIS_TOPIC,AUTOSTART,MQTT_BROKER_HOST,MQTT_BROKER_PORT,MQTT_CLIENT_ID,MQTT_PASSWORD,MQTT_USERNAME,OUTPUT_QOS,RISK_TOPIC,ROOT_CONFIDENCE,SOURCE_QOS,SOURCE_TOPIC,datetime,json,math,mo,mqtt,os,pd,threading,time,timezone,utc_now_iso,) -
单元 2 - 定义 Bayesian Network 拓扑、状态空间、CPT 生成器和概率规则。
显示代码
Terminal window @app.celldef _(math, pd):"""Static Bayesian-network definition and CPT generators."""import itertoolsfrom pgmpy.factors.discrete import TabularCPDfrom pgmpy.inference import VariableEliminationfrom pgmpy.models import DiscreteBayesianNetworkSTATE_SPACES = {"STATE_3": ["LOW", "NORMAL", "HIGH"],"STATE_YN": ["NO", "YES"],"PB2_3STATE": ["FAIL", "IL_OK", "SL_OK"],"TOP_EVENT_STATE": ["NORMAL", "DEVELOPING", "CONFIRMED"],}NODE_STATE_SPACE = {"T1_HighChlorideLoad": "STATE_3","T2_ChemProtectionFail": "STATE_3","T3_WashWaterLow": "STATE_3","PB1_Desalter_OK": "STATE_YN","PB2_WashQty_State": "PB2_3STATE","PB3_Neutralization_OK": "STATE_YN","MB1_ESD_OK": "STATE_YN","MB2_TSV_OK": "STATE_YN","TE_CorrosionState": "TOP_EVENT_STATE","EV_TOTIRON_High": "STATE_YN","EV_SAP_Positive": "STATE_YN","EV_QuillFailure": "STATE_YN","C1_LOPC": "STATE_YN","C2_UnplannedShutdown": "STATE_YN","C3_Escalation": "STATE_YN",}EDGES = [("T1_HighChlorideLoad", "TE_CorrosionState"),("T2_ChemProtectionFail", "TE_CorrosionState"),("T3_WashWaterLow", "TE_CorrosionState"),("PB1_Desalter_OK", "TE_CorrosionState"),("PB2_WashQty_State", "TE_CorrosionState"),("PB3_Neutralization_OK", "TE_CorrosionState"),("TE_CorrosionState", "EV_TOTIRON_High"),("TE_CorrosionState", "EV_SAP_Positive"),("TE_CorrosionState", "EV_QuillFailure"),("TE_CorrosionState", "C1_LOPC"),("MB1_ESD_OK", "C1_LOPC"),("MB2_TSV_OK", "C1_LOPC"),("TE_CorrosionState", "C2_UnplannedShutdown"),("MB1_ESD_OK", "C2_UnplannedShutdown"),("C1_LOPC", "C3_Escalation"),("MB1_ESD_OK", "C3_Escalation"),("MB2_TSV_OK", "C3_Escalation"),]def states_for(node_id: str) -> list[str]:return STATE_SPACES[NODE_STATE_SPACE[node_id]]def softmax(values: list[float]) -> list[float]:maximum = max(values)exps = [math.exp(value - maximum) for value in values]total = sum(exps)return [value / total for value in exps]THREAT_SCORE = {"LOW": 0.0, "NORMAL": 1.0, "HIGH": 2.0}PB1_OK_BONUS = {"NO": 0.8, "YES": -0.6}PB3_OK_BONUS = {"NO": 1.0, "YES": -0.7}PB2_STATE_BONUS = {"FAIL": 1.2, "IL_OK": 0.3, "SL_OK": -0.6}def top_event_distribution(t1, t2, t3, pb1, pb2, pb3) -> dict[str, float]:score = (THREAT_SCORE[t1]+ THREAT_SCORE[t2]+ THREAT_SCORE[t3]+ PB1_OK_BONUS[pb1]+ PB2_STATE_BONUS[pb2]+ PB3_OK_BONUS[pb3])probabilities = softmax([2.2 - score,-0.2 + 0.6 * score,-1.2 + 0.8 * score,])return dict(zip(STATE_SPACES["TOP_EVENT_STATE"], probabilities))TE_PARENTS = ["T1_HighChlorideLoad","T2_ChemProtectionFail","T3_WashWaterLow","PB1_Desalter_OK","PB2_WashQty_State","PB3_Neutralization_OK",]TE_PARENT_SPACES = [states_for(parent) for parent in TE_PARENTS]te_rows = []for combination in itertools.product(*TE_PARENT_SPACES):distribution = top_event_distribution(*combination)te_rows.append([*combination,distribution["NORMAL"],distribution["DEVELOPING"],distribution["CONFIRMED"],])TE_CPT_DF = pd.DataFrame(te_rows,columns=TE_PARENTS+ ["P_NORMAL", "P_DEVELOPING", "P_CONFIRMED"],)def evidence_yes_probability(te_state: str, kind: str) -> float:mappings = {"TOTIRON": {"NORMAL": 0.15, "DEVELOPING": 0.65, "CONFIRMED": 0.85},"SAP": {"NORMAL": 0.08, "DEVELOPING": 0.35, "CONFIRMED": 0.75},"QUILL": {"NORMAL": 0.03, "DEVELOPING": 0.15, "CONFIRMED": 0.55},}return mappings[kind][te_state]def lopc_yes_probability(te: str, mb1: str, mb2: str) -> float:base = {"NORMAL": 0.02, "DEVELOPING": 0.08, "CONFIRMED": 0.25}[te]return min(0.98, base + (0.10 if mb1 == "NO" else 0) + (0.18 if mb2 == "NO" else 0))def shutdown_yes_probability(te: str, mb1: str) -> float:base = {"NORMAL": 0.05, "DEVELOPING": 0.35, "CONFIRMED": 0.60}[te]return min(0.98, base + (0.10 if mb1 == "NO" else 0))def escalation_yes_probability(lopc: str, mb1: str, mb2: str) -> float:if lopc == "NO":return 0.01return min(0.98, 0.10 + (0.18 if mb1 == "NO" else 0) + (0.35 if mb2 == "NO" else 0))def peaked_prior(node_id: str, center_state: str, confidence: float) -> list[float]:states = states_for(node_id)confidence = min(max(float(confidence), 0.0), 1.0)if center_state not in states:return [1.0 / len(states)] * len(states)if len(states) == 1:return [1.0]remainder = (1.0 - confidence) / (len(states) - 1)probabilities = [remainder] * len(states)probabilities[states.index(center_state)] = confidencereturn probabilitiesdef binary_cpd(child: str, parent: str, yes_probability_fn):parent_states = states_for(parent)yes_values = [float(yes_probability_fn(state)) for state in parent_states]return TabularCPD(variable=child,variable_card=2,values=[[1.0 - value for value in yes_values], yes_values],evidence=[parent],evidence_card=[len(parent_states)],state_names={child: states_for(child), parent: parent_states},)def binary_multi_cpd(child: str, parents: list[str], yes_probability_fn):parent_spaces = [states_for(parent) for parent in parents]combinations = list(itertools.product(*parent_spaces))yes_values = [float(yes_probability_fn(*combination)) for combination in combinations]return TabularCPD(variable=child,variable_card=2,values=[[1.0 - value for value in yes_values], yes_values],evidence=parents,evidence_card=[len(space) for space in parent_spaces],state_names={child: states_for(child),**{parent: states_for(parent) for parent in parents},},)return (DiscreteBayesianNetwork,EDGES,NODE_STATE_SPACE,STATE_SPACES,TE_CPT_DF,TE_PARENTS,TabularCPD,VariableElimination,binary_cpd,binary_multi_cpd,escalation_yes_probability,evidence_yes_probability,itertools,lopc_yes_probability,peaked_prior,shutdown_yes_probability,states_for,) -
单元 3 - 校验传入的 MQTT 消息内容,并将连续过程值转换为 Bayesian Network 状态。
显示代码
Terminal window @app.celldef _(ROOT_CONFIDENCE, utc_now_iso):"""Source validation, discretization, and one-snapshot analysis."""import hashlibREQUIRED_SOURCE_FIELDS = ["d103_chloride","d103_ph","wash_water_flow","desalter_salt_ptb","desalter_bsw","wash_water_rate","rrd_ph","rrd_chloride","rrd_total_iron","sap_evidence","injection_quill_fail","timestamp",]def parse_boolean(value) -> bool:if isinstance(value, bool):return valueif isinstance(value, (int, float)):return bool(value)if isinstance(value, str):normalized = value.strip().lower()if normalized in {"true", "1", "yes", "y", "on"}:return Trueif normalized in {"false", "0", "no", "n", "off"}:return Falseraise ValueError(f"invalid boolean value: {value!r}")def normalize_sap(value) -> str:normalized = str(value).strip().lower()if normalized in {"0", "none", "normal", "", "nan"}:return "None"if normalized in {"1", "work_order", "work order", "minor", "warning"}:return "Minor corrosion observed"if normalized in {"2","inspection_finding","inspection finding","confirmed","abnormal",}:return "Confirmed corrosion / thickness loss"raise ValueError(f"unsupported sap_evidence value: {value!r}")def normalize_source_payload(payload: dict) -> dict:if not isinstance(payload, dict):raise TypeError("source MQTT payload must be a JSON object")lower = {str(key).lower(): value for key, value in payload.items()}missing = [field for field in REQUIRED_SOURCE_FIELDS if field not in lower]if missing:raise ValueError(f"source payload is missing fields: {missing}")normalized = {"d103_chloride": float(lower["d103_chloride"]),"d103_ph": float(lower["d103_ph"]),"wash_water_flow": float(lower["wash_water_flow"]),"desalter_salt_ptb": float(lower["desalter_salt_ptb"]),"desalter_bsw": float(lower["desalter_bsw"]),"wash_water_rate": float(lower["wash_water_rate"]),"rrd_ph": float(lower["rrd_ph"]),"rrd_chloride": float(lower["rrd_chloride"]),"rrd_total_iron": float(lower["rrd_total_iron"]),"sap_evidence": normalize_sap(lower["sap_evidence"]),"injection_quill_fail": parse_boolean(lower["injection_quill_fail"]),"timestamp": str(lower["timestamp"]).strip(),}numeric_ranges = {"d103_chloride": (0.0, 1000.0),"d103_ph": (0.0, 14.0),"wash_water_flow": (0.0, 10000.0),"desalter_salt_ptb": (0.0, 1000.0),"desalter_bsw": (0.0, 100.0),"wash_water_rate": (0.0, 100.0),"rrd_ph": (0.0, 14.0),"rrd_chloride": (0.0, 10000.0),"rrd_total_iron": (0.0, 10000.0),}for field, (minimum, maximum) in numeric_ranges.items():value = normalized[field]if not minimum <= value <= maximum:raise ValueError(f"{field} out of accepted range: {value}")if not normalized["timestamp"]:raise ValueError("timestamp cannot be empty")return normalizeddef source_message_id(payload: dict) -> str:# Timestamp is the primary event key. Hash protects against two different# snapshots accidentally sharing the same timestamp.ordered = "|".join(str(payload[field]) for field in REQUIRED_SOURCE_FIELDS)return hashlib.sha256(ordered.encode("utf-8")).hexdigest()def discretize_source(payload: dict) -> tuple[dict, dict]:t1 = ("LOW"if payload["d103_chloride"] < 5.0else "NORMAL"if payload["d103_chloride"] <= 20.0else "HIGH")chemical_failure_score = int(payload["d103_ph"] < 5.5) + int(payload["d103_chloride"] > 10.0)t2 = ("HIGH"if chemical_failure_score == 2else "NORMAL"if chemical_failure_score == 1else "LOW")t3 = ("HIGH"if payload["wash_water_flow"] < 6.0else "NORMAL"if payload["wash_water_flow"] <= 10.0else "LOW")pb1 = ("YES"if payload["desalter_salt_ptb"] <= 1.0and payload["desalter_bsw"] <= 0.2else "NO")pb2 = ("FAIL"if payload["wash_water_rate"] < 4.0else "IL_OK"if payload["wash_water_rate"] < 6.0else "SL_OK")pb3 = ("YES"if payload["rrd_ph"] >= 5.5 and payload["rrd_chloride"] <= 10.0else "NO")root_centers = {"T1_HighChlorideLoad": t1,"T2_ChemProtectionFail": t2,"T3_WashWaterLow": t3,"PB1_Desalter_OK": pb1,"PB2_WashQty_State": pb2,"PB3_Neutralization_OK": pb3,# The source model currently has no MB fields. Keep both mitigations# available by default, matching the previous notebook UI defaults."MB1_ESD_OK": "YES","MB2_TSV_OK": "YES",}evidence = {"EV_TOTIRON_High": "YES" if payload["rrd_total_iron"] >= 10.0 else "NO","EV_SAP_Positive": "YES" if payload["sap_evidence"] != "None" else "NO","EV_QuillFailure": "YES" if payload["injection_quill_fail"] else "NO",}return root_centers, evidencereturn (REQUIRED_SOURCE_FIELDS,discretize_source,normalize_source_payload,parse_boolean,source_message_id,) -
单元 4 - 基于当前快照构建 Bayesian Network,执行推理并生成输出消息内容。
显示代码
Terminal window @app.celldef _(DiscreteBayesianNetwork,EDGES,ROOT_CONFIDENCE,STATE_SPACES,TE_CPT_DF,TE_PARENTS,TabularCPD,VariableElimination,binary_cpd,binary_multi_cpd,discretize_source,escalation_yes_probability,evidence_yes_probability,itertools,lopc_yes_probability,peaked_prior,shutdown_yes_probability,states_for,utc_now_iso,):"""Build a fresh BN for the current snapshot and execute inference."""def factor_probability(factor, variable: str, state: str) -> float:states = factor.state_names[variable]values = list(map(float, factor.values))return float(values[states.index(state)])def build_model(root_centers: dict):model = DiscreteBayesianNetwork(EDGES)root_cpds = []for node_id, center_state in root_centers.items():states = states_for(node_id)probabilities = peaked_prior(node_id, center_state, ROOT_CONFIDENCE)root_cpds.append(TabularCPD(variable=node_id,variable_card=len(states),values=[[probability] for probability in probabilities],state_names={node_id: states},))te_values = [TE_CPT_DF["P_NORMAL"].astype(float).tolist(),TE_CPT_DF["P_DEVELOPING"].astype(float).tolist(),TE_CPT_DF["P_CONFIRMED"].astype(float).tolist(),]te_cpd = TabularCPD(variable="TE_CorrosionState",variable_card=3,values=te_values,evidence=TE_PARENTS,evidence_card=[len(states_for(parent)) for parent in TE_PARENTS],state_names={"TE_CorrosionState": STATE_SPACES["TOP_EVENT_STATE"],**{parent: states_for(parent) for parent in TE_PARENTS},},)evidence_cpds = [binary_cpd("EV_TOTIRON_High","TE_CorrosionState",lambda state: evidence_yes_probability(state, "TOTIRON"),),binary_cpd("EV_SAP_Positive","TE_CorrosionState",lambda state: evidence_yes_probability(state, "SAP"),),binary_cpd("EV_QuillFailure","TE_CorrosionState",lambda state: evidence_yes_probability(state, "QUILL"),),]consequence_cpds = [binary_multi_cpd("C1_LOPC",["TE_CorrosionState", "MB1_ESD_OK", "MB2_TSV_OK"],lopc_yes_probability,),binary_multi_cpd("C2_UnplannedShutdown",["TE_CorrosionState", "MB1_ESD_OK"],shutdown_yes_probability,),binary_multi_cpd("C3_Escalation",["C1_LOPC", "MB1_ESD_OK", "MB2_TSV_OK"],escalation_yes_probability,),]model.add_cpds(*root_cpds, te_cpd, *evidence_cpds, *consequence_cpds)if not model.check_model():raise RuntimeError("Bayesian network model validation failed")return modeldef analyze_snapshot(payload: dict, previous_risk_state: str | None = None) -> dict:root_centers, evidence = discretize_source(payload)model = build_model(root_centers)inference = VariableElimination(model)q_te = inference.query(["TE_CorrosionState"], evidence=evidence, show_progress=False)q_c1 = inference.query(["C1_LOPC"], evidence=evidence, show_progress=False)q_c2 = inference.query(["C2_UnplannedShutdown"], evidence=evidence, show_progress=False)q_c3 = inference.query(["C3_Escalation"], evidence=evidence, show_progress=False)p_normal = factor_probability(q_te, "TE_CorrosionState", "NORMAL")p_developing = factor_probability(q_te, "TE_CorrosionState", "DEVELOPING")p_confirmed = factor_probability(q_te, "TE_CorrosionState", "CONFIRMED")probabilities = {"NORMAL": p_normal,"DEVELOPING": p_developing,"CONFIRMED": p_confirmed,}risk_state = max(probabilities, key=probabilities.get)analysis_timestamp = utc_now_iso()analysis_payload = {"p_normal": p_normal,"p_developing": p_developing,"p_confirmed": p_confirmed,"total_corrosion_risk": p_developing + p_confirmed,"lopc_probability": factor_probability(q_c1, "C1_LOPC", "YES"),"shutdown_probability": factor_probability(q_c2, "C2_UnplannedShutdown", "YES"),"escalation_probability": factor_probability(q_c3, "C3_Escalation", "YES"),"timestamp": analysis_timestamp,}risk_payload = {"risk_state": risk_state,"previous_state": previous_risk_state or risk_state,"model_confidence": max(probabilities.values()),"timestamp": analysis_timestamp,}return {"source_timestamp": payload["timestamp"],"root_states": root_centers,"evidence": evidence,"analysis_payload": analysis_payload,"risk_payload": risk_payload,"risk_state": risk_state,}return analyze_snapshot, build_model, factor_probability -
单元 5 - 订阅源 topic,处理每个传入快照,执行 Bayesian 分析并发布结果。
显示代码
Terminal window @app.celldef _(ANALYSIS_TOPIC,AUTOSTART,MQTT_BROKER_HOST,MQTT_BROKER_PORT,MQTT_CLIENT_ID,MQTT_PASSWORD,MQTT_USERNAME,OUTPUT_QOS,RISK_TOPIC,SOURCE_QOS,SOURCE_TOPIC,analyze_snapshot,json,mqtt,normalize_source_payload,source_message_id,threading,utc_now_iso,):"""Long-running MQTT subscriber and publisher service."""class CorrosionBNService:def __init__(self):self.lock = threading.RLock()self.client = Noneself.running = Falseself.connected = Falseself.last_message_id = Noneself.last_source_timestamp = Noneself.last_analysis_timestamp = Noneself.previous_risk_state = Noneself.processed_count = 0self.duplicate_count = 0self.error_count = 0self.last_error = Noneself.last_result = Noneself.connection_rc = Nonedef _create_client(self):# Use MQTT 3.1.1 and callback API v1 for broad compatibility with# EMQX and different paho-mqtt releases bundled in Notebook images.kwargs = {"client_id": MQTT_CLIENT_ID,"clean_session": True,"protocol": mqtt.MQTTv311,"transport": "tcp",}try:client = mqtt.Client(callback_api_version=mqtt.CallbackAPIVersion.VERSION1,**kwargs,)except (AttributeError, TypeError):client = mqtt.Client(**kwargs)client.username_pw_set(MQTT_USERNAME, MQTT_PASSWORD)client.on_connect = self._on_connectclient.on_disconnect = self._on_disconnectclient.on_message = self._on_messageclient.reconnect_delay_set(min_delay=1, max_delay=30)client.enable_logger()return clientdef _on_connect(self, client, userdata, flags, rc, properties=None):code = int(rc)with self.lock:self.connected = code == 0self.connection_rc = codeself.last_error = None if code == 0 else f"MQTT connection rejected, rc={code}"if code == 0:sub_rc, mid = client.subscribe(SOURCE_TOPIC, qos=SOURCE_QOS)if sub_rc != mqtt.MQTT_ERR_SUCCESS:with self.lock:self.last_error = f"MQTT subscribe failed, rc={sub_rc}"returnprint(f"[BN SERVICE] connected and subscribed to {SOURCE_TOPIC}")else:print(f"[BN SERVICE ERROR] connection rejected rc={code}")def _on_disconnect(self, client, userdata, rc, properties=None):with self.lock:self.connected = Falseself.connection_rc = int(rc) if rc is not None else Noneif self.running and rc:self.last_error = f"MQTT disconnected unexpectedly, rc={rc}"print(f"[BN SERVICE] disconnected: {rc}")def _publish_json(self, topic: str, payload: dict):info = self.client.publish(topic,payload=json.dumps(payload, separators=(",", ":"), ensure_ascii=False),qos=OUTPUT_QOS,retain=False,)info.wait_for_publish(timeout=10)if info.rc != mqtt.MQTT_ERR_SUCCESS:raise RuntimeError(f"MQTT publish failed for {topic}: rc={info.rc}")def _on_message(self, client, userdata, msg):try:raw_payload = json.loads(msg.payload.decode("utf-8"))payload = normalize_source_payload(raw_payload)message_id = source_message_id(payload)with self.lock:if message_id == self.last_message_id:self.duplicate_count += 1returnprevious_risk_state = self.previous_risk_stateresult = analyze_snapshot(payload, previous_risk_state)# Publish only after both payloads have been generated successfully.self._publish_json(ANALYSIS_TOPIC, result["analysis_payload"])self._publish_json(RISK_TOPIC, result["risk_payload"])with self.lock:self.last_message_id = message_idself.last_source_timestamp = payload["timestamp"]self.last_analysis_timestamp = result["analysis_payload"]["timestamp"]self.previous_risk_state = result["risk_state"]self.processed_count += 1self.last_result = resultself.last_error = Noneprint("[BN SERVICE] processed",payload["timestamp"],"->",result["risk_state"],)except Exception as exc:with self.lock:self.error_count += 1self.last_error = f"{type(exc).__name__}: {exc}"print("[BN SERVICE ERROR]", self.last_error)def start(self):with self.lock:if self.running:return Falseif not MQTT_PASSWORD:raise RuntimeError("TIER0_MQTT_PASSWORD is empty. Set it before starting the service.")self.client = self._create_client()self.running = Truetry:# Synchronous connect surfaces DNS, socket, and authentication# failures immediately instead of leaving running=True forever.rc = self.client.connect(MQTT_BROKER_HOST,MQTT_BROKER_PORT,keepalive=30,)if rc != mqtt.MQTT_ERR_SUCCESS:raise RuntimeError(f"MQTT connect() failed, rc={rc}")self.client.loop_start()return Trueexcept Exception as exc:with self.lock:self.running = Falseself.connected = Falseself.last_error = f"{type(exc).__name__}: {exc}"self.client = Noneraisedef stop(self):with self.lock:if not self.running:return Falseclient = self.clientself.running = Falseself.connected = Falseself.client = Nonetry:client.disconnect()finally:client.loop_stop()return Truedef restart(self):self.stop()return self.start()def status(self) -> dict:with self.lock:return {"running": self.running,"connected": self.connected,"connection_rc": self.connection_rc,"client_id": MQTT_CLIENT_ID,"broker": f"{MQTT_BROKER_HOST}:{MQTT_BROKER_PORT}","source_topic": SOURCE_TOPIC,"analysis_topic": ANALYSIS_TOPIC,"risk_topic": RISK_TOPIC,"processed_count": self.processed_count,"duplicate_count": self.duplicate_count,"error_count": self.error_count,"last_source_timestamp": self.last_source_timestamp,"last_analysis_timestamp": self.last_analysis_timestamp,"previous_risk_state": self.previous_risk_state,"last_error": self.last_error,}# Reuse a previous service instance when the cell is re-executed, preventing# duplicate subscribers inside the same Python process.existing = globals().get("_CORROSION_BN_SERVICE")if existing is not None:try:existing.stop()except Exception:passCORROSION_BN_SERVICE = CorrosionBNService()globals()["_CORROSION_BN_SERVICE"] = CORROSION_BN_SERVICEAUTOSTART_ERROR = Noneif AUTOSTART:try:CORROSION_BN_SERVICE.start()except Exception as exc:AUTOSTART_ERROR = f"{type(exc).__name__}: {exc}"return AUTOSTART_ERROR, CORROSION_BN_SERVICE, CorrosionBNService -
单元 6 - 显示 MQTT 连接状态和服务统计信息。
显示代码
Terminal window @app.celldef _(AUTOSTART_ERROR, CORROSION_BN_SERVICE, mo, pd):"""Service status and usage instructions."""status = CORROSION_BN_SERVICE.status()blocks = [mo.md("## MQTT analysis service"),pd.DataFrame([status]),]if status.get("last_error"):blocks.append(mo.md(f"### Connection / processing error\n`{status['last_error']}`"))elif status.get("running") and not status.get("connected"):blocks.append(mo.md("### Connecting\nThe MQTT client has started but has not completed the broker connection yet."))if AUTOSTART_ERROR:blocks.append(mo.md(f"""### Service not started`{AUTOSTART_ERROR}`Set the MQTT password in the Notebook environment and run this cell again:~~~pythonimport os# MQTT credentials are already configured in this notebook.~~~"""))else:blocks.append(mo.md("""The service starts automatically and processes one complete sourcesnapshot for every new MQTT message. No manual analysis trigger or SQLpolling is required.Available controls in Python:~~~pythonCORROSION_BN_SERVICE.status()CORROSION_BN_SERVICE.stop()CORROSION_BN_SERVICE.start()CORROSION_BN_SERVICE.restart()~~~"""))mo.vstack(blocks)return status, -
单元 7 - 显示最新的 Bayesian 分析结果。
显示代码
Terminal window @app.celldef _(CORROSION_BN_SERVICE, mo):"""Display the most recent analysis result when one is available."""latest = CORROSION_BN_SERVICE.last_resultif latest is None:mo.md("### Latest result\nNo source message has been processed yet.")else:mo.vstack([mo.md("### Latest source snapshot result"),mo.md(f"- Source timestamp: `{latest['source_timestamp']}`"),mo.md(f"- Risk state: `{latest['risk_state']}`"),mo.md("#### Root states"),latest["root_states"],mo.md("#### Evidence"),latest["evidence"],mo.md("#### Corrosion analysis payload"),latest["analysis_payload"],mo.md("#### Corrosion risk payload"),latest["risk_payload"],])returnif __name__ == "__main__":app.run()
-
-
运行所有单元,然后进入 UNS 检查结果。
构建 Bow-tie 应用
Section titled “构建 Bow-tie 应用”-
在 Tier0 中进入 Builder。
-
在对话框中输入应用需求,并开始构建。
显示 Builder prompt
Terminal window 为 CDU 常压塔顶腐蚀构建一个现代工业 Bow-Tie Analysis 应用。仅使用以下 UNS topic 作为数据源:- Refinery/CDU_Plant/Metric/Corrosion_Monitoring- Refinery/CDU_Plant/Metric/Corrosion_Analysis- Refinery/CDU_Plant/State/Corrosion_Risk当新的 UNS 消息到达时,应用必须自动更新。该应用应帮助操作员理解当前腐蚀风险、相关威胁、屏障状态以及可能后果。创建一个单页 dashboard,包含:1. 将 Bow-Tie 图作为主区域:- 左侧显示威胁- 在威胁与顶事件之间显示预防性屏障- 以腐蚀作为中心顶事件- 在顶事件之后显示缓解性屏障- 右侧显示后果2. 显著展示风险摘要,包括:- 当前风险状态- 模型置信度- 正常、发展中、已确认三类概率- 总腐蚀风险3. 一个紧凑的实时过程数据面板,显示 Corrosion_Monitoring 中的关键值,包括 chloride、pH、wash water flow、desalter performance、wash water rate、total iron、SAP evidence 和 injection quill status。4. 后果指标包括:- 一次包容失效- 非计划停机- 升级风险5. 视觉状态规则:- 绿色表示正常或健康状态- 琥珀色/黄色表示风险发展中或警告状态- 红色表示风险已确认、屏障失效或严重后果- 灰色表示数据不可用- 蓝色仅用于交互和信息高亮6. 交互:- 选择威胁、屏障或后果时,应显示其当前状态、相关过程值和简短说明。- 当新的 UNS 数据到达时,自动更新界面。- 显示最新数据时间戳,以及内置 SDK 提供的连接/订阅状态。设计要求:在整个应用中保持现有视觉设计语言和样式。使用现代工业运营 dashboard 风格:- 使用深蓝灰色调,而不是纯黑色。- 使用带有细微渐变的分层面板。- 使用卡片式布局,圆角较大,约 16-20px。- 使用细边框做分隔,避免厚重阴影。- 使用柔和层次感和克制的发光效果。- 保持清晰间距、充足内边距和一致对齐。- 正文使用 IBM Plex Sans。- 标题使用 Space Grotesk。- 使用紧凑、信息密集且层级清晰的卡片。- 使用增加字距的大写 section label。- 呈现专业的炼化厂 / SCADA 界面,而不是商业分析 dashboard。- 使用扁平插图和 SVG 图形,避免拟物元素。- 使用平滑、细微的过渡和 hover 效果。整体外观应清晰、现代、技术化,并适合运营场景。避免:- 纯黑背景- 玻璃拟态- 赛博朋克风格- 霓虹灯效- 过度装饰效果- 商业分析 dashboard 风格使用清晰、专业的炼化运营风格,保持明确层级、紧凑卡片、可读标签和响应式布局。严格使用对应 UNS 模型中定义的 payload 字段。不要虚构额外字段、计算值或 mock 数据。 -
经过几轮对话,确认应用完成后,部署应用。
-
进入 应用中心,打开应用并检查。