Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions bkmonitor/config/default.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"

Expand Down
Original file line number Diff line number Diff line change
@@ -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="物化视图状态"),
),
]
4 changes: 4 additions & 0 deletions bkmonitor/metadata/models/data_link/data_link_configs.py
Original file line number Diff line number Diff line change
Expand Up @@ -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绑定配置"
Expand Down
184 changes: 184 additions & 0 deletions bkmonitor/metadata/service/surrealdb_materialized_view.py
Original file line number Diff line number Diff line change
@@ -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
31 changes: 30 additions & 1 deletion bkmonitor/metadata/task/bkbase.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand Down
104 changes: 104 additions & 0 deletions bkmonitor/metadata/tests/service/test_surrealdb_materialized_view.py
Original file line number Diff line number Diff line change
@@ -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"
Loading
Loading