"""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_domain"
# FQDN はインターネット向けの接尾辞で終わるものに限定する。
# これにより、part-000.csv のようなパスをドメインとして誤検出しない。
FQDN_RE = re.compile(
r"(?<![A-Za-z0-9_.-])"
r"((?:[A-Za-z0-9](?:[A-Za-z0-9-]{0,61}[A-Za-z0-9])?\.)+"
r"(?:com|net|org|io|co|ai|dev|app|edu|gov|jp|uk|de|fr|au|ca|us))"
r"(?=$|[\s/:?&=,'\"<>;)\]}.])",
re.IGNORECASE,
)
# URL にだけ現れる operation-service:5501 のような短縮ホストも抽出する。
# 任意の単語は対象にせず、実際の URL の場合だけを対象とする。
URL_HOST_RE = re.compile(
r"(?i)\b(?:https?|abfss?)://(?:[^/@\s]+@)?"
r"([a-z0-9](?:[a-z0-9.-]*[a-z0-9])?)(?::\d+)?"
)
S3_RE = re.compile(r"(?i)\bs3://([a-z0-9][a-z0-9.-]{1,61}[a-z0-9])(?:/[^\s'\"<>]*)?")
LOCAL_OR_CLASS = re.compile(
r"(?i)^(?:localhost|ip6-localhost|ip6-loopback|"
r"(?:java|javax|org|com|net|scala|sun)\.(?:io|net|util|lang))$"
)
# Java/Scala のパッケージ名は FQDN に似るため除外する。
# ただし、URL_HOST_RE で取得した実際の URI ホストは除外しない。
CLASS_PACKAGE_PREFIX_RE = re.compile(
r"(?i)^(?:com\.sun|org\.apache|org\.sparkproject|java|javax|scala|sun)(?:\.|$)"
)
def normalized_host(host: str) -> str:
"""可変のホスト名をワイルドカード化し、接続先の系統を保持する。"""
h = host.lower().rstrip(".")
# spark2triprodje、spark3triprodje、GUID 付きの派生ホストを同じ系統にまとめる。
m = re.search(r"(?:^|\.)spark\d+triprodje\.dfs\.core\.windows\.net$", h)
if m:
return "*.spark*triprodje.dfs.core.windows.net"
# Fabric は OneLake の前にワークスペース/容量固有のプレフィックスを付与する。
if h.endswith(".onelake.fabric.microsoft.com"):
return "*.onelake.fabric.microsoft.com"
# 数字はサービス種別ではなくインスタンス識別子である。
if re.match(r"^tokenservice\d+\.japaneast\.trident\.azuresynapse\.net$", h):
return "tokenservice*.japaneast.trident.azuresynapse.net"
# ノートブックリソースのマウントに使われる Fabric の管理ストレージホスト。
# 生成された実ホスト名は raw_target_evidence に保持する。
if re.match(r"^olst[a-z0-9]+\.dfs\.core\.windows\.net$", h):
return "olst*.dfs.core.windows.net"
return h
def target_scope(target_type: str, domain: str) -> str:
"""Microsoft が公開している範囲だけを内部扱いとし、不明な名前は確認対象に残す。"""
if target_type == "S3 URI":
return "Fabric 外部(確認済み)"
# Microsoft が公開している Fabric/OneLake エンドポイントだけを内部として扱う。
if domain.endswith(".fabric.microsoft.com"):
return "Fabric(Microsoft 公開)"
# これらは Fabric 実行時ログで確認されるが、Microsoft による所有範囲の公開情報がない。
# 内部と断定せず、確認対象として保持する。
if (
domain == "operation-service"
or domain.endswith(".notebook.windows.net")
or domain.endswith(".trident.azuresynapse.net")
or domain.endswith(".pbidedicated.windows.net")
or domain == "*.spark*triprodje.dfs.core.windows.net"
or domain == "olst*.dfs.core.windows.net"
):
return "Fabric 内部の可能性(要確認)"
# 汎用 Azure Storage エンドポイントは利用者の ADLS アカウントである可能性がある。
# 隣接する API/成功ログと併せて確認する必要がある。
return "要確認"
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_log_targets(text: str):
"""1 行のログから、種別・正規化後の対象・元の対象を返す。"""
if not text:
return []
out = set()
for m in S3_RE.finditer(text):
bucket = m.group(1).lower()
# S3 バケットは DNS ホストではないため、スキームを含めて保持する。
out.add(("S3 URI", f"s3://{bucket}", m.group(0)))
for m in URL_HOST_RE.finditer(text):
raw = m.group(1).lower()
if raw == "operation-service":
out.add(("内部サービス URL", raw, raw))
elif "." in raw:
out.add(("ドメイン", normalized_host(raw), raw))
for m in FQDN_RE.finditer(text):
raw = m.group(1).lower()
trailing_text = text[m.end():m.end() + 2]
if (
raw.endswith((".local", ".internal"))
or LOCAL_OR_CLASS.fullmatch(raw)
or CLASS_PACKAGE_PREFIX_RE.match(raw)
# 文末のピリオドは許可するが、パッケージ名の一部は許可しない。
or (
trailing_text.startswith(".")
and len(trailing_text) > 1
and re.match(r"[A-Za-z0-9_-]", trailing_text[1])
)
or re.match(r"^127(?:\.\d{1,3}){3}$", raw)
or raw == "0.0.0.0"
):
continue
out.add(("ドメイン", normalized_host(raw), raw))
return [Row(target_type=x[0], domain=x[1], raw_target=x[2]) for x in sorted(out)]
target_schema = T.ArrayType(
T.StructType([
T.StructField("target_type", T.StringType(), False),
T.StructField("domain", T.StringType(), False),
T.StructField("raw_target", T.StringType(), False),
])
)
extract_targets_udf = F.udf(extract_log_targets, target_schema)
target_scope_udf = F.udf(target_scope, T.StringType())
collapse_file_names_udf = F.udf(collapse_file_names, T.StringType())
source = spark.table(SOURCE_TABLE)
# 1 行に複数のドメイン/URI が存在し得るため、先に展開する。
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("target", F.explode_outer(extract_targets_udf("log_text")))
.where(F.col("target").isNotNull())
.select(
"workspace_id", "notebook_id", "livy_id", "spark_application_id",
"submitted_datetime", "file_name", "line_number",
F.col("target.target_type").alias("target_type"),
F.col("target.domain").alias("log_line_context_logtext_domain"),
F.col("target.raw_target").alias("raw_target_evidence"),
)
.withColumn("target_scope", target_scope_udf("target_type", "log_line_context_logtext_domain"))
)
# まず同じファイル内の重複をまとめ、その後ノートブック実行内の全ファイルで同一対象をまとめる。
execution_keys = [
"workspace_id", "notebook_id", "livy_id", "spark_application_id",
"submitted_datetime", "target_type", "target_scope",
"log_line_context_logtext_domain",
]
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_target_evidence")).alias("raw_target_evidence"),
)
)
# 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_target_evidence")))).alias("raw_target_evidence"),
)
.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_domain", "target_type", "target_scope",
"file_name_members", "line_number_by_file", "raw_target_evidence",
)
)
(result.write.format("delta").mode("overwrite").option("overwriteSchema", "true").saveAsTable(TARGET_TABLE))
display(result.orderBy("submitted_datetime", "file_name", "log_line_context_logtext_domain"))
detect_domain
