"""Fabric のドライバーログからクライアント/コネクターを抽出する。"""
from pyspark.sql import Row, functions as F, types as T
import re
SOURCE_TABLE = "input_allp_audit_trn_spark_driver_log_line"
TARGET_TABLE = "input_allp_audit_trn_spark_driver_log_line_service_client"
# スタックトレースと実行時メッセージからクライアントクラスを取得する。
# FooServiceClientBuilder と FooServiceClient は同じクライアントとして正規化する。
CLIENT_CLASS_RE = re.compile(
r"\b([A-Za-z][A-Za-z0-9]*(?:ServiceClient(?:Builder)?|ClientImpl|S3Client|HttpClient))\b"
)
# コネクターとデータベースドライバーのクラスを取得する。
CONNECTOR_CLASS_RE = re.compile(
r"\b([A-Za-z][A-Za-z0-9]*(?:Connector[A-Za-z0-9]*|Driver))\b"
)
# JAR/WHL 名はクラスパス、パッケージ解決、例外メッセージに現れる。
# connector/driver を含む名前に限定し、無関係なライブラリを検出しない。
CONNECTOR_ARTIFACT_RE = re.compile(
r"(?i)\b([A-Za-z0-9][A-Za-z0-9._+-]*(?:connector|driver)[A-Za-z0-9._+-]*\.(?:jar|whl))\b"
)
# パッケージの存在は依存関係が利用可能なことを示すだけで、API 呼出しや書込みの証拠ではない。
PACKAGE_CLIENT_RULES = [
(re.compile(r"(?i)\bsnowflake-connector-python\b"), "SnowflakeConnectorPython"),
(re.compile(r"(?i)\bboto3\b"), "AWSBoto3"),
(re.compile(r"(?i)\bazure-storage-file-datalake\b"), "AzureStorageDataLakeSDK"),
(re.compile(r"(?i)\bazure-storage-blob\b"), "AzureStorageBlobSDK"),
(re.compile(r"(?i)\boffice365-rest-python-client\b"), "SharePointOffice365SDK"),
(re.compile(r"(?i)\bshareplum\b"), "SharePlum"),
(re.compile(r"(?i)\bmsgraph-sdk\b"), "MicrosoftGraphSDK"),
(re.compile(r"(?i)\bpyodbc\b"), "PyODBC"),
(re.compile(r"(?i)\bpymssql\b"), "PyMSSQL"),
(re.compile(r"(?i)\bpsycopg(?:2)?\b"), "PsycopgPostgreSQL"),
(re.compile(r"(?i)\bpymongo\b"), "PyMongo"),
(re.compile(r"(?i)\bgoogle-cloud-storage\b"), "GoogleCloudStorageSDK"),
(re.compile(r"(?i)\bgoogle-cloud-bigquery\b"), "GoogleCloudBigQuerySDK"),
]
# これらは Fabric ランタイムの実装要素であり、ノートブックの業務処理で使うサービスクライアントではない。
EXCLUDED_COMPONENTS = {"InMemoryCacheClient", "Connector", "Driver"}
def normalize_client_name(name: str) -> str:
name = re.sub(r"(?:Builder|Impl)$", "", name)
return name
def normalize_connector_artifact(name: str) -> str:
"""拡張子と末尾のバージョンを除去し、製品名を残す。"""
base = re.sub(r"(?i)\.(?:jar|whl)$", "", name)
return re.sub(r"-(?:v)?\d.*$", "", base)
def collapse_file_names(file_names):
"""同一対象のファイル名を、必要最小限の * パターンで表現する。"""
names = sorted({x for x in (file_names or []) if x})
if not names:
return None
if len(names) == 1:
return names[0]
prefix = names[0]
for name in names[1:]:
i = 0
while i < min(len(prefix), len(name)) and prefix[i] == name[i]:
i += 1
prefix = prefix[:i]
reversed_names = [name[::-1] for name in names]
suffix = reversed_names[0]
for name in reversed_names[1:]:
i = 0
while i < min(len(suffix), len(name)) and suffix[i] == name[i]:
i += 1
suffix = suffix[:i]
suffix = suffix[::-1]
if len(prefix) + len(suffix) >= min(map(len, names)):
suffix = ""
return f"{prefix}*{suffix}"
def extract_service_clients(text: str):
"""種別と証拠レベルを含むコンポーネント一覧を返す。"""
if not text:
return []
found = set()
for match in CLIENT_CLASS_RE.finditer(text):
raw = match.group(1)
normalized = normalize_client_name(raw)
if normalized not in EXCLUDED_COMPONENTS:
found.add(("サービスクライアントクラス", "実行時参照", normalized, raw))
for match in CONNECTOR_CLASS_RE.finditer(text):
raw = match.group(1)
if raw not in EXCLUDED_COMPONENTS:
found.add(("コネクター/ドライバークラス", "実行時参照", raw, raw))
for match in CONNECTOR_ARTIFACT_RE.finditer(text):
raw = match.group(1)
found.add(("コネクター/ドライバー成果物", "依存関係のみ", normalize_connector_artifact(raw), raw))
for pattern, client_name in PACKAGE_CLIENT_RULES:
if pattern.search(text):
found.add(("SDK パッケージ", "依存関係のみ", client_name, client_name))
return [
Row(
service_client_signal_kind=x[0],
component_evidence_level=x[1],
service_client=x[2],
raw_service_client_evidence=x[3],
)
for x in sorted(found)
]
service_client_schema = T.ArrayType(
T.StructType([
T.StructField("service_client_signal_kind", T.StringType(), False),
T.StructField("component_evidence_level", T.StringType(), False),
T.StructField("service_client", T.StringType(), False),
T.StructField("raw_service_client_evidence", T.StringType(), False),
])
)
extract_service_clients_udf = F.udf(extract_service_clients, service_client_schema)
collapse_file_names_udf = F.udf(collapse_file_names, T.StringType())
source = spark.table(SOURCE_TABLE)
detected = (
source
.select(
"workspace_id", "notebook_id", "livy_id", "spark_application_id",
"submitted_datetime", "file_name", "line_number",
F.col("log_line_context_logtext").alias("log_text"),
)
.withColumn("service_client_signal", F.explode_outer(extract_service_clients_udf("log_text")))
.where(F.col("service_client_signal").isNotNull())
.select(
"workspace_id", "notebook_id", "livy_id", "spark_application_id",
"submitted_datetime", "file_name", "line_number",
F.col("service_client_signal.service_client").alias("log_line_context_logtext_service_client"),
F.col("service_client_signal.service_client_signal_kind").alias("service_client_signal_kind"),
F.col("service_client_signal.component_evidence_level").alias("component_evidence_level"),
F.col("service_client_signal.raw_service_client_evidence").alias("raw_service_client_evidence"),
F.when(
F.col("log_text").rlike(r"(?i)\b(?:warn(?:ing)?|error|fatal|exception)\b"),
F.col("log_text"),
).alias("service_client_alert_log"),
)
)
# まずファイル内の重複をまとめ、その後、同一ノートブック実行の全ログ種別をまとめる。
execution_keys = [
"workspace_id", "notebook_id", "livy_id", "spark_application_id",
"submitted_datetime", "service_client_signal_kind", "component_evidence_level",
"log_line_context_logtext_service_client",
]
per_file = (
detected
.groupBy(*execution_keys, "file_name")
.agg(
F.sort_array(F.collect_set(F.col("line_number").cast("bigint"))).alias("line_number"),
F.sort_array(F.collect_set("raw_service_client_evidence")).alias("raw_service_client_evidence"),
F.sort_array(F.collect_set("service_client_alert_log")).alias("service_client_alert_log"),
)
)
# file_name をワイルドカード化した後も、元のファイル名と行番号の対応を保持する。
result = (
per_file
.groupBy(*execution_keys)
.agg(
F.sort_array(F.collect_set("file_name")).alias("file_name_members"),
F.sort_array(F.array_distinct(F.flatten(F.collect_list("line_number")))).alias("line_number"),
F.map_from_entries(F.collect_list(F.struct("file_name", "line_number"))).alias("line_number_by_file"),
F.sort_array(F.array_distinct(F.flatten(F.collect_list("raw_service_client_evidence")))).alias("raw_service_client_evidence"),
F.array_distinct(F.flatten(F.collect_list("service_client_alert_log"))).alias("service_client_alert_log"),
)
.withColumn("file_name", collapse_file_names_udf("file_name_members"))
.select(
"workspace_id", "notebook_id", "livy_id", "spark_application_id",
"submitted_datetime", "file_name", "line_number",
"log_line_context_logtext_service_client", "service_client_signal_kind",
"component_evidence_level",
"file_name_members", "line_number_by_file", "raw_service_client_evidence",
"service_client_alert_log",
)
)
(result.write.format("delta").mode("overwrite").option("overwriteSchema", "true").saveAsTable(TARGET_TABLE))
display(
result.orderBy(
"submitted_datetime", "file_name",
"log_line_context_logtext_service_client",
)
)
detect_serviceclient
