From 04005955ca5bd93ce7cb5dd764e98200e64326d9 Mon Sep 17 00:00:00 2001 From: Zhaoyikaiii Date: Mon, 17 Aug 2026 17:47:33 +0800 Subject: [PATCH] =?UTF-8?q?feat(metadata):=20=E4=B8=8B=E5=8F=91=20SurrealD?= =?UTF-8?q?B=20=E7=89=A9=E5=8C=96=E8=A7=86=E5=9B=BE=E5=AE=9A=E4=B9=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- bkmonitor/config/default.py | 3 + .../0276_surrealdb_materialized_view_state.py | 30 +++ .../models/data_link/data_link_configs.py | 4 + .../service/surrealdb_materialized_view.py | 184 ++++++++++++++++++ bkmonitor/metadata/task/bkbase.py | 31 ++- .../test_surrealdb_materialized_view.py | 104 ++++++++++ .../tests/task/test_refresh_data_link.py | 39 ++++ 7 files changed, 394 insertions(+), 1 deletion(-) create mode 100644 bkmonitor/metadata/migrations/0276_surrealdb_materialized_view_state.py create mode 100644 bkmonitor/metadata/service/surrealdb_materialized_view.py create mode 100644 bkmonitor/metadata/tests/service/test_surrealdb_materialized_view.py diff --git a/bkmonitor/config/default.py b/bkmonitor/config/default.py index 0746d0cc2be..d7bd5ce4e99 100644 --- a/bkmonitor/config/default.py +++ b/bkmonitor/config/default.py @@ -1668,6 +1668,9 @@ def decode_plaintext(cls, plaintext_bytes: bytes, encoding: str = "utf-8", **kwa # 是否开启空间内置数据链路初始化 ENABLE_SPACE_BUILTIN_DATA_LINK = os.getenv("ENABLE_SPACE_BUILTIN_DATA_LINK", "false").lower() == "true" +# 是否在 SurrealDBBinding 就绪后下发 relation materialized view DDL,默认关闭 +ENABLE_SURREALDB_MATERIALIZED_VIEW = os.getenv("BKAPP_ENABLE_SURREALDB_MATERIALIZED_VIEW", "false").lower() == "true" + # 创建 vm 链路资源所属的命名空间 DEFAULT_VM_DATA_LINK_NAMESPACE = "bkmonitor" diff --git a/bkmonitor/metadata/migrations/0276_surrealdb_materialized_view_state.py b/bkmonitor/metadata/migrations/0276_surrealdb_materialized_view_state.py new file mode 100644 index 00000000000..a1fef62bd3c --- /dev/null +++ b/bkmonitor/metadata/migrations/0276_surrealdb_materialized_view_state.py @@ -0,0 +1,30 @@ +from django.db import migrations, models + + +class Migration(migrations.Migration): + dependencies = [ + ("metadata", "0275_bcsfederalclusterinfo_tenant_scope"), + ] + + operations = [ + migrations.AddField( + model_name="surrealdbbindingconfig", + name="materialized_view_definition_hash", + field=models.CharField(default="", max_length=64, verbose_name="物化视图定义哈希"), + ), + migrations.AddField( + model_name="surrealdbbindingconfig", + name="materialized_view_last_apply_time", + field=models.DateTimeField(blank=True, null=True, verbose_name="物化视图最近下发时间"), + ), + migrations.AddField( + model_name="surrealdbbindingconfig", + name="materialized_view_last_error", + field=models.TextField(blank=True, default="", verbose_name="物化视图最近错误"), + ), + migrations.AddField( + model_name="surrealdbbindingconfig", + name="materialized_view_status", + field=models.CharField(blank=True, default="", max_length=64, verbose_name="物化视图状态"), + ), + ] diff --git a/bkmonitor/metadata/models/data_link/data_link_configs.py b/bkmonitor/metadata/models/data_link/data_link_configs.py index 383b1a7c6f2..2c257e2da17 100644 --- a/bkmonitor/metadata/models/data_link/data_link_configs.py +++ b/bkmonitor/metadata/models/data_link/data_link_configs.py @@ -1110,6 +1110,10 @@ class SurrealDBBindingConfig(DataLinkResourceConfigBase): table_type = models.CharField(verbose_name="图表类型", max_length=32, default="temporary") vertices = models.JSONField(verbose_name="顶点定义", default=list) relations = models.JSONField(verbose_name="关系定义", default=list) + materialized_view_definition_hash = models.CharField(verbose_name="物化视图定义哈希", max_length=64, default="") + materialized_view_status = models.CharField(verbose_name="物化视图状态", max_length=64, default="", blank=True) + materialized_view_last_error = models.TextField(verbose_name="物化视图最近错误", default="", blank=True) + materialized_view_last_apply_time = models.DateTimeField(verbose_name="物化视图最近下发时间", null=True, blank=True) class Meta: verbose_name = "SurrealDB绑定配置" diff --git a/bkmonitor/metadata/service/surrealdb_materialized_view.py b/bkmonitor/metadata/service/surrealdb_materialized_view.py new file mode 100644 index 00000000000..a1c446c1482 --- /dev/null +++ b/bkmonitor/metadata/service/surrealdb_materialized_view.py @@ -0,0 +1,184 @@ +"""SurrealDB relation materialized view definition and reconciliation.""" + +from __future__ import annotations + +import hashlib +import json +import re +from dataclasses import dataclass +from typing import TYPE_CHECKING, Any + +if TYPE_CHECKING: + from metadata.models.data_link.data_link_configs import SurrealDBBindingConfig + +_IDENTIFIER_PATTERN = re.compile(r"^[A-Za-z_][A-Za-z0-9_]*$") +_SCOPE_ANNOTATION_KEYS = { + "namespace": ("surrealdbnamespace", "surrealnamespace"), + "database": ("surrealdbdatabase", "surrealdatabase"), +} + + +@dataclass(frozen=True) +class SurrealDBScope: + namespace: str + database: str + + +def _normalize_key(value: str) -> str: + return value.replace("_", "").replace("-", "").lower() + + +def _require_identifier(value: Any, field: str) -> str: + if not isinstance(value, str) or not _IDENTIFIER_PATTERN.fullmatch(value): + raise ValueError(f"{field} 不是合法的 SurrealDB 标识符: {value!r}") + return value + + +def _quote_identifier(value: str) -> str: + return f"`{value}`" + + +def _string_literal(value: str) -> str: + return json.dumps(value, ensure_ascii=False) + + +def resolve_surrealdb_scope(config: dict[str, Any]) -> SurrealDBScope: + metadata = config.get("metadata") if isinstance(config, dict) else None + annotations = metadata.get("annotations") if isinstance(metadata, dict) else None + if not isinstance(annotations, dict): + annotations = {} + normalized_annotations = {_normalize_key(key): value for key, value in annotations.items() if isinstance(key, str)} + + resolved = {} + for field, keys in _SCOPE_ANNOTATION_KEYS.items(): + value = next((normalized_annotations.get(key) for key in keys if normalized_annotations.get(key)), None) + if value is not None: + resolved[field] = _require_identifier(value, f"metadata.annotations.{field}") + + status = config.get("status") if isinstance(config, dict) else None + scope = status.get("storage") if isinstance(status, dict) else None + if isinstance(scope, dict): + for field in ("namespace", "database"): + if field not in resolved and scope.get(field): + resolved[field] = _require_identifier(scope[field], f"status.storage.{field}") + + if set(resolved) != {"namespace", "database"}: + raise ValueError("SurrealDBBinding 缺少 SurrealDB namespace/database annotations") + return SurrealDBScope(namespace=resolved["namespace"], database=resolved["database"]) + + +def _snapshot_expression(record_link: str, fields: list[str], field: str) -> str: + projections = [] + for index, item in enumerate(fields): + identifier = _require_identifier(item, f"{field}[{index}]") + projections.append(f" {identifier}: {record_link}.{identifier}") + return "{\n" + ",\n".join(projections) + "\n }" + + +def build_materialized_view_ddl(binding: SurrealDBBindingConfig, scope: SurrealDBScope) -> str: + binding._validate_graph_definitions() + vertices = {vertex["name"]: vertex for vertex in binding.vertices} + statements = [ + f"USE NS {_quote_identifier(scope.namespace)} DB {_quote_identifier(scope.database)};", + ] + + for index, relation in enumerate(binding.relations): + relation_name = _require_identifier(relation["name"], f"relations[{index}].name") + source_type = _require_identifier(relation["from"], f"relations[{index}].from") + target_type = _require_identifier(relation["to"], f"relations[{index}].to") + source_vertex = vertices.get(source_type) + target_vertex = vertices.get(target_type) + if source_vertex is None or target_vertex is None: + raise ValueError(f"relation[{relation_name}] 引用了未定义的顶点") + + view_name = f"{relation_name}_materialized_view" + source_table = f"{relation_name}_liveness_record" + source_index = f"idx_{relation_name}_mv_source_period" + target_index = f"idx_{relation_name}_mv_target_period" + source_snapshot = _snapshot_expression( + "relation_id.in", source_vertex["id_fields"], f"vertices[{source_type}].id_fields" + ) + target_snapshot = _snapshot_expression( + "relation_id.out", target_vertex["id_fields"], f"vertices[{target_type}].id_fields" + ) + + statements.extend( + [ + f"REMOVE TABLE IF EXISTS {_quote_identifier(view_name)};", + "\n".join( + [ + f"DEFINE TABLE {_quote_identifier(view_name)} TYPE NORMAL AS", + "SELECT", + " id AS liveness_id,", + " relation_id,", + " relation_id.in AS source_id,", + " relation_id.out AS target_id,", + f" {_string_literal(source_type)} AS source_type,", + f" {_string_literal(target_type)} AS target_type,", + f" {source_snapshot} AS source_snapshot,", + f" {target_snapshot} AS target_snapshot,", + " period_start AS relation_period_start_ms,", + " period_end AS relation_period_end_ms", + f"FROM {_quote_identifier(source_table)}", + "WHERE period_start < period_end;", + ] + ), + "\n".join( + [ + f"DEFINE INDEX {_quote_identifier(source_index)}", + f"ON TABLE {_quote_identifier(view_name)}", + "FIELDS source_id, relation_period_start_ms, relation_period_end_ms;", + ] + ), + "\n".join( + [ + f"DEFINE INDEX {_quote_identifier(target_index)}", + f"ON TABLE {_quote_identifier(view_name)}", + "FIELDS target_id, relation_period_start_ms, relation_period_end_ms;", + ] + ), + ] + ) + return "\n\n".join(statements) + + +def reconcile_materialized_views(binding: SurrealDBBindingConfig, remote_config: dict[str, Any]) -> bool: + from django.utils import timezone + + from core.drf_resource import api + from metadata.models.data_link.constants import DataLinkResourceStatus + + if binding.status != DataLinkResourceStatus.OK.value: + return False + + scope = resolve_surrealdb_scope(remote_config) + ddl = build_materialized_view_ddl(binding, scope) + definition_hash = hashlib.sha256(ddl.encode("utf-8")).hexdigest() + if ( + binding.materialized_view_status == DataLinkResourceStatus.OK.value + and binding.materialized_view_definition_hash == definition_hash + ): + return False + + try: + api.bkdata.query_data(sql=ddl, prefer_storage="surrealdb") + except Exception as error: + binding.materialized_view_status = DataLinkResourceStatus.FAILED.value + binding.materialized_view_last_error = str(error)[:4096] + binding.save(update_fields=["materialized_view_status", "materialized_view_last_error", "last_modify_time"]) + raise + + binding.materialized_view_definition_hash = definition_hash + binding.materialized_view_status = DataLinkResourceStatus.OK.value + binding.materialized_view_last_error = "" + binding.materialized_view_last_apply_time = timezone.now() + binding.save( + update_fields=[ + "materialized_view_definition_hash", + "materialized_view_status", + "materialized_view_last_error", + "materialized_view_last_apply_time", + "last_modify_time", + ] + ) + return True diff --git a/bkmonitor/metadata/task/bkbase.py b/bkmonitor/metadata/task/bkbase.py index b6d9a6ebb74..fb10a4b8e0e 100644 --- a/bkmonitor/metadata/task/bkbase.py +++ b/bkmonitor/metadata/task/bkbase.py @@ -39,8 +39,14 @@ DataLinkKind, DataLinkResourceStatus, ) -from metadata.models.data_link.data_link_configs import COMPONENT_CLASS_MAP, ClusterConfig, ResultTableConfig +from metadata.models.data_link.data_link_configs import ( + COMPONENT_CLASS_MAP, + ClusterConfig, + ResultTableConfig, + SurrealDBBindingConfig, +) from metadata.models.space.constants import SpaceStatus, SpaceTypes +from metadata.service.surrealdb_materialized_view import reconcile_materialized_views from metadata.models.vm.utils import report_metadata_data_link_status_info from metadata.service.sync_metadata import sync_kafka_metadata, sync_vm_metadata from metadata.task.constants import BKBASE_V4_KIND_STORAGE_CONFIGS @@ -1472,6 +1478,8 @@ def _reconcile_data_link_components() -> tuple[ stats.untrusted_batch_count += 1 continue + remote_configs_by_name = {config["metadata"]["name"]: config for config in configs} + if component_kind in STORAGE_BINDING_KIND_MAP: reference_issues = _check_storage_binding_references( configs, @@ -1544,6 +1552,27 @@ def _reconcile_data_link_components() -> tuple[ stats.untrusted_batch_count += 1 continue + if component_kind == DataLinkKind.SURREALDBBINDING.value and settings.ENABLE_SURREALDB_MATERIALIZED_VIEW: + materialized_view_components = { + component.name: component for component in [*components, *created_components] + } + for name, (_, extra_config) in parsed_configs.items(): + component = materialized_view_components.get(name) + if not isinstance(component, SurrealDBBindingConfig): + continue + component.status = extra_config["status"] + try: + reconcile_materialized_views(component, remote_configs_by_name[name]) + except Exception as error: # pylint: disable=broad-except + logger.exception( + "bulk_refresh_data_link_status: reconcile surrealdb materialized views failed, " + "tenant->[%s], namespace->[%s], name->[%s], error->[%s]", + bk_tenant_id, + namespace, + name, + error, + ) + stats.created_count += len(created_components) stats.updated_count += len(changed_components) - terminated_count stats.terminated_count += terminated_count diff --git a/bkmonitor/metadata/tests/service/test_surrealdb_materialized_view.py b/bkmonitor/metadata/tests/service/test_surrealdb_materialized_view.py new file mode 100644 index 00000000000..72c771c8896 --- /dev/null +++ b/bkmonitor/metadata/tests/service/test_surrealdb_materialized_view.py @@ -0,0 +1,104 @@ +from unittest.mock import ANY + +import pytest + +from metadata.models.data_link.constants import DataLinkResourceStatus +from metadata.models.data_link.data_link_configs import SurrealDBBindingConfig +from metadata.service.surrealdb_materialized_view import ( + SurrealDBScope, + build_materialized_view_ddl, + reconcile_materialized_views, + resolve_surrealdb_scope, +) + + +def _binding(**kwargs): + defaults = { + "name": "graph_binding", + "bk_biz_id": 2, + "status": DataLinkResourceStatus.OK.value, + "vertices": [ + {"name": "node", "id_fields": ["bcs_cluster_id", "node"]}, + {"name": "pod", "id_fields": ["bcs_cluster_id", "namespace", "pod"]}, + ], + "relations": [{"name": "node_with_pod", "from": "node", "to": "pod"}], + } + defaults.update(kwargs) + return SurrealDBBindingConfig(**defaults) + + +def _remote_config(): + return { + "metadata": { + "annotations": { + "SurrealDBNamespace": "bkmonitor", + "surrealdb_database": "biz_2", + } + }, + "status": {"phase": DataLinkResourceStatus.OK.value}, + } + + +def test_resolve_surrealdb_scope_from_annotations(): + assert resolve_surrealdb_scope(_remote_config()) == SurrealDBScope(namespace="bkmonitor", database="biz_2") + + +def test_resolve_surrealdb_scope_rejects_missing_database(): + with pytest.raises(ValueError, match="namespace/database"): + resolve_surrealdb_scope({"metadata": {"annotations": {"SurrealDBNamespace": "bkmonitor"}}}) + + +def test_build_materialized_view_ddl(): + ddl = build_materialized_view_ddl(_binding(), SurrealDBScope(namespace="bkmonitor", database="biz_2")) + + assert "USE NS `bkmonitor` DB `biz_2`;" in ddl + assert "REMOVE TABLE IF EXISTS `node_with_pod_materialized_view`;" in ddl + assert "DEFINE TABLE `node_with_pod_materialized_view` TYPE NORMAL AS" in ddl + assert "relation_id.in AS source_id" in ddl + assert "node: relation_id.in.node" in ddl + assert "pod: relation_id.out.pod" in ddl + assert "period_end AS relation_period_end_ms\nFROM `node_with_pod_liveness_record`" in ddl + assert "FIELDS source_id, relation_period_start_ms, relation_period_end_ms;" in ddl + assert "FIELDS target_id, relation_period_start_ms, relation_period_end_ms;" in ddl + + +def test_build_materialized_view_ddl_rejects_unsafe_identifier(): + with pytest.raises(ValueError, match="合法的 SurrealDB 标识符"): + build_materialized_view_ddl( + _binding(relations=[{"name": "node;REMOVE", "from": "node", "to": "pod"}]), + SurrealDBScope(namespace="bkmonitor", database="biz_2"), + ) + + +@pytest.mark.django_db(databases="__all__") +def test_reconcile_materialized_views_applies_once(mocker): + binding = _binding() + binding.save() + query_data = mocker.patch("core.drf_resource.api.bkdata.query_data") + + assert reconcile_materialized_views(binding, _remote_config()) is True + binding.refresh_from_db() + assert binding.materialized_view_status == DataLinkResourceStatus.OK.value + assert len(binding.materialized_view_definition_hash) == 64 + assert binding.materialized_view_last_apply_time is not None + query_data.assert_called_once_with(sql=ANY, prefer_storage="surrealdb") + + assert reconcile_materialized_views(binding, _remote_config()) is False + query_data.assert_called_once() + + +@pytest.mark.django_db(databases="__all__") +def test_reconcile_materialized_views_records_failure(mocker): + binding = _binding() + binding.save() + mocker.patch( + "core.drf_resource.api.bkdata.query_data", + side_effect=RuntimeError("query failed"), + ) + + with pytest.raises(RuntimeError, match="query failed"): + reconcile_materialized_views(binding, _remote_config()) + + binding.refresh_from_db() + assert binding.materialized_view_status == DataLinkResourceStatus.FAILED.value + assert binding.materialized_view_last_error == "query failed" diff --git a/bkmonitor/metadata/tests/task/test_refresh_data_link.py b/bkmonitor/metadata/tests/task/test_refresh_data_link.py index 213126e0c64..b7d67aa9921 100644 --- a/bkmonitor/metadata/tests/task/test_refresh_data_link.py +++ b/bkmonitor/metadata/tests/task/test_refresh_data_link.py @@ -855,6 +855,45 @@ def test_refresh_updates_empty_surrealdb_definitions(mocker): assert component.relations == [] +@pytest.mark.django_db(databases="__all__") +def test_refresh_reconciles_surrealdb_materialized_views_when_enabled(mocker, settings): + settings.ENABLE_SURREALDB_MATERIALIZED_VIEW = True + component = models.SurrealDBBindingConfig.objects.create( + name="graph_rt", + namespace="bkmonitor", + bk_tenant_id="system", + data_link_name="graph_link", + bk_biz_id=2, + status=DataLinkResourceStatus.PENDING.value, + surrealdb_cluster_name="surreal-default", + bkbase_result_table_name="graph_rt", + table_type="normal", + vertices=[{"name": "pod", "id_fields": ["pod_name"]}], + relations=[{"name": "pod_with_pod", "from": "pod", "to": "pod"}], + ) + remote_config = _remote_component_for_kind(DataLinkKind.SURREALDBBINDING.value, "graph_rt") + remote_config["spec"].update( + { + "storage": {"name": "surreal-default"}, + "data": {"name": "graph_rt"}, + "vertices": component.vertices, + "relations": component.relations, + } + ) + remote_config["metadata"]["annotations"] = { + "SurrealDBNamespace": "bkmonitor", + "SurrealDBDatabase": "biz_2", + } + mocker.patch("metadata.task.bkbase.api.bkdata.list_data_link", return_value=[remote_config]) + reconcile = mocker.patch("metadata.task.bkbase.reconcile_materialized_views") + + _reconcile_data_link_components() + + component.refresh_from_db() + assert component.status == DataLinkResourceStatus.OK.value + reconcile.assert_called_once_with(component, remote_config) + + @pytest.mark.django_db(databases="__all__") def test_refresh_keeps_falsy_non_surrealdb_fields(mocker): component = models.DataIdConfig.objects.create(