FROM DATA TO DIRECTION

BIVIET

Đối Tác Tin Cậy

Hơn 15 năm kinh nghiệm BI tại Nhật Bản
Đối tác tin cậy cho doanh nghiệp Việt và Nhật
日本で15年以上のBI導入経験
日本とベトナムの企業を支える、信頼のデータパートナー
detect_domain
0
(0)
"""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"))

How useful was this post?

Click on a star to rate it!

Average rating 0 / 5. Vote count: 0

No votes so far! Be the first to rate this post.

detect_domain

Để lại một bình luận

Email của bạn sẽ không được hiển thị công khai. Các trường bắt buộc được đánh dấu *

Chuyển lên trên