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_result_before_exclution
0
(0)
"""ドライバーログから、外部データ送信の根拠を集約して確認対象を作成する。"""

from pyspark.sql import functions as F

SOURCE_LOG_TABLE = "input_allp_audit_trn_spark_driver_log_line"
DOMAIN_TABLE = "input_allp_audit_trn_spark_driver_log_line_domain"
SERVICE_CLIENT_TABLE = "input_allp_audit_trn_spark_driver_log_line_service_client"
TARGET_TABLE = "input_allp_audit_trn_spark_driver_log_line_detected_result"

GROUP_KEYS = ["workspace_id", "notebook_id", "livy_id", "submitted_datetime"]

# write や save のような一般語は Fabric/Spark 内部の書込みにも現れるため使わない。
# 外部送信を示す具体的な API 名、成功メッセージ、HTTP メソッドだけを対象にする。
ADLS_BLOB_WRITE_RE = r"(?i)\b(?:append_data|flush_data|upload_data|upload_blob|stage_block|commit_block_list|create_file)\b"
S3_WRITE_RE = r"(?i)\b(?:putobject|put_object|upload_file|upload_fileobj|complete_multipart_upload)\b"
SNOWFLAKE_WRITE_RE = r"(?i)\b(?:copy\s+into|insert\s+into|write_pandas|inserted\s+\d+\s+rows?\s+into)\b"
SHAREPOINT_WRITE_RE = r"(?i)\b(?:createuploadsession|files/add|add_using_path|uploadsession)\b"
HTTP_WRITE_RE = r"(?i)\b(?:http(?:s)?\s+(?:request|method).{0,80}\b(?:put|post|patch)\b|(?:put|post|patch)\s+https?://)"

domain_source = spark.table(DOMAIN_TABLE)
service_client_source = spark.table(SERVICE_CLIENT_TABLE)
source_log = spark.table(SOURCE_LOG_TABLE)
service_client_columns = set(service_client_source.columns)

def component_column(name, default_value=None):
    """旧スキーマのテーブルでも実行できるよう、不足列に既定値を設定する。"""
    return F.col(name) if name in service_client_columns else F.lit(default_value)

# ドメイン抽出テーブルは、監査証跡としてログ中の endpoint をすべて保持する。
# この検知結果だけで、Fabric の公開バックエンドと確認済みの内部/一時領域を除外する。
# 汎用 Azure Storage の *.dfs.core.windows.net、*.blob.core.windows.net は除外しない。
FABRIC_INTERNAL_ENDPOINT_RE = (
    r"(?i)^(?:"
    # 「fabric」を含む endpoint は Fabric 内部 endpoint として検知対象から除外する。
    r".*fabric.*|"
    # 確認済みの Fabric Spark 実行時サービス/一時ストレージ。
    r"operation-service|"
    r"\*\.spark\*triprodje\.dfs\.core\.windows\.net|"
    r"(?:[a-z0-9-]+\.)*spark\d+triprodje\.dfs\.core\.windows\.net|"
    r"(?:[a-z0-9*-]+\.)*pbidedicated\.windows\.net|"
    r"exec\.japaneast\.notebook\.windows\.net|"
    r"olst\*\.dfs\.core\.windows\.net|"
    r"olst[a-z0-9]+\.dfs\.core\.windows\.net|"
    r"tokenservice\*\.japaneast\.trident\.azuresynapse\.net|"
    r"tokenservice\d+\.japaneast\.trident\.azuresynapse\.net|"
    r"rpjapaneast\.svc\.datafactory\.azure\.com|"
    # Microsoft が公開している Fabric バックエンド endpoint。
    r"(?:[a-z0-9*-]+\.)*fabric\.microsoft\.com|"
    r"(?:[a-z0-9*-]+\.)*pbidedicated\.microsoft\.com|"
    r"(?:[a-z0-9*-]+\.)*analysis\.windows\.net|"
    r"(?:[a-z0-9*-]+\.)*powerbi\.com|"
    r"content\.powerapps\.com|"
    r"res\.cdn\.office\.net|"
    r"aznb-ame-prod\.azureedge\.net|"
    r"aznbcdn\.notebooks\.azure\.net|"
    r"(?:[a-z0-9*-]+\.)*notebooks\.azuresandbox\.ms"
    r")$"
)
AZURE_STORAGE_ACCOUNT_ENDPOINT_RE = r"(?i)^(?:[a-z0-9-]+\.)+(?:dfs|blob)\.core\.windows\.net$"
SNOWFLAKE_ENDPOINT_RE = r"(?i)^(?:[a-z0-9*-]+\.)*snowflakecomputing\.com$"
EXTERNAL_DESTINATION_KINDS = ["Amazon S3", "Snowflake", "Azure Storage Account"]
INTERNAL_FABRIC_RUNTIME_SERVICE_CLIENTS = ["ApacheHttpClient", "BlobServiceClient", "TokenServiceClient"]
INTERNAL_FABRIC_COMPONENT_RE = r"(?i)(?:spark|fabric)"
# 外部送信先の確定に使わない既知の SDK/コネクター/ドライバー。
# 必要に応じて、確認済みの名称をこの一覧に追加する。
EXCLUDED_CONNECTOR_OR_DRIVER_COMPONENTS = [
    "AzureStorageBlobSDK",
    "VPaaSConnector",
    "VegasConnector",
    "gcs-connector-hadoop3",
    "mysql-connector-java",
]

# 除外判定では NULL を false として扱い、S3 URI や Snowflake を誤って落とさない。
domain_prepared = (
    domain_source
    .withColumn(
        "is_fabric_internal_endpoint",
        F.coalesce(F.col("log_line_context_logtext_domain").rlike(FABRIC_INTERNAL_ENDPOINT_RE), F.lit(False)),
    )
    .withColumn(
        "destination_kind",
        F.when(F.col("target_type") == "S3 URI", F.lit("Amazon S3"))
         .when(F.col("log_line_context_logtext_domain").rlike(SNOWFLAKE_ENDPOINT_RE), F.lit("Snowflake"))
         .when(F.col("log_line_context_logtext_domain").rlike(AZURE_STORAGE_ACCOUNT_ENDPOINT_RE), F.lit("Azure Storage Account"))
         .otherwise(F.lit("ドメイン確認対象")),
    )
)

# S3 URI は bucket 名に "fabric" を含み得るため、内部 endpoint の文字列判定から常に除外する。
# これにより s3://... の検知結果を誤って落とさない。
non_fabric_domain = domain_prepared.where(
    (F.col("target_type") == "S3 URI") | ~F.col("is_fabric_internal_endpoint")
)

domain_evidence_json = F.to_json(
    F.struct(
        F.col("spark_application_id"),
        F.col("log_line_context_logtext_domain").alias("target"),
        F.col("target_type"), F.col("target_scope"), F.col("destination_kind"), F.col("file_name"),
        F.col("file_name_members"), F.col("line_number_by_file"),
        F.col("raw_target_evidence"),
    )
)
domain_group = (
    non_fabric_domain.withColumn("domain_evidence", domain_evidence_json)
    .groupBy(*GROUP_KEYS)
    .agg(
        F.sort_array(F.collect_set("spark_application_id")).alias("domain_spark_application_ids"),
        F.sort_array(F.collect_set("log_line_context_logtext_domain")).alias("non_fabric_targets"),
        F.sort_array(F.collect_set("target_scope")).alias("target_scopes"),
        F.sort_array(F.collect_set("destination_kind")).alias("domain_destination_kinds"),
        F.sort_array(F.collect_set("domain_evidence")).alias("non_fabric_domain_evidence"),
    )
)

component_prepared = service_client_source.select(
    *GROUP_KEYS, F.col("spark_application_id"),
    F.col("log_line_context_logtext_service_client").alias("component_name"),
    component_column("service_client_signal_kind", "unknown").alias("component_signal_kind"),
    component_column("component_evidence_level", "unknown").alias("component_evidence_level"),
    F.col("file_name"), F.col("file_name_members"), F.col("line_number_by_file"),
    F.col("raw_service_client_evidence"),
    component_column("service_client_alert_log").alias("service_client_alert_log"),
).where(
    # BlobServiceClient と TokenServiceClient は Fabric ランタイムの内部処理として扱い、
    # ノートブックから外部サービスへ送信するクライアントとしては集計しない。
    ~F.col("log_line_context_logtext_service_client").isin(INTERNAL_FABRIC_RUNTIME_SERVICE_CLIENTS)
    & ~F.col("log_line_context_logtext_service_client").isin(EXCLUDED_CONNECTOR_OR_DRIVER_COMPONENTS)
    & ~F.col("log_line_context_logtext_service_client").rlike(INTERNAL_FABRIC_COMPONENT_RE)
)
component_evidence_json = F.to_json(
    F.struct(
        F.col("spark_application_id"), F.col("component_name"),
        F.col("component_signal_kind"), F.col("component_evidence_level"),
        F.col("file_name"), F.col("file_name_members"), F.col("line_number_by_file"),
        F.col("raw_service_client_evidence"), F.col("service_client_alert_log"),
    )
)
component_group = (
    component_prepared.withColumn("component_evidence", component_evidence_json)
    .groupBy(*GROUP_KEYS)
    .agg(
        F.sort_array(F.collect_set("spark_application_id")).alias("component_spark_application_ids"),
        F.sort_array(F.collect_set("component_name")).alias("components_observed"),
        F.sort_array(F.collect_set(F.when(
            F.col("component_signal_kind") == "サービスクライアントクラス", F.col("component_name")
        ))).alias("runtime_service_clients"),
        F.sort_array(F.collect_set(F.when(
            F.col("component_signal_kind") == "コネクター/ドライバークラス", F.col("component_name")
        ))).alias("runtime_connectors_or_drivers"),
        F.sort_array(F.collect_set(F.when(
            F.col("component_signal_kind").isin("コネクター/ドライバー成果物", "SDK パッケージ"), F.col("component_name")
        ))).alias("connector_or_driver_dependencies"),
        F.sort_array(F.collect_set("component_signal_kind")).alias("component_signal_kinds"),
        F.sort_array(F.collect_set("component_evidence")).alias("component_evidence"),
    )
)

# 具体的な送信操作のログだけを、ファイル名・行番号付きで保持する。
# 同じノートブック実行内で発生した根拠であり、対象との対応は evidence で確認する。
action_prepared = (
    source_log.select(
        *GROUP_KEYS, "spark_application_id", "file_name", "line_number",
        F.col("log_line_context_logtext").alias("log_text"),
    )
    .withColumn(
        "outbound_action_kind",
        F.when(F.col("log_text").rlike(ADLS_BLOB_WRITE_RE), F.lit("書込み操作"))
         .when(F.col("log_text").rlike(S3_WRITE_RE), F.lit("書込み操作"))
         .when(F.col("log_text").rlike(SNOWFLAKE_WRITE_RE), F.lit("書込み操作"))
         .when(F.col("log_text").rlike(SHAREPOINT_WRITE_RE), F.lit("書込み操作"))
         .when(F.col("log_text").rlike(HTTP_WRITE_RE), F.lit("書込み操作")),
    )
    .where(F.col("outbound_action_kind").isNotNull())
)
outbound_action_evidence_json = F.to_json(
    F.struct(
        F.col("spark_application_id"), F.col("outbound_action_kind"),
        F.col("file_name"), F.col("line_number").cast("bigint").alias("line_number"),
        F.col("log_text"),
    )
)
action_group = (
    action_prepared.withColumn("outbound_action_evidence_item", outbound_action_evidence_json)
    .groupBy(*GROUP_KEYS)
    .agg(
        F.sort_array(F.collect_set("spark_application_id")).alias("action_spark_application_ids"),
        F.sort_array(F.collect_set("outbound_action_kind")).alias("outbound_action_kinds"),
        F.sort_array(F.collect_set("outbound_action_evidence_item")).alias("outbound_action_evidence"),
    )
)

empty_string_array = F.array().cast("array<string>")
base = (
    domain_group.join(component_group, GROUP_KEYS, "full_outer")
    .join(action_group, GROUP_KEYS, "full_outer")
    .withColumn(
        "spark_application_ids",
        F.sort_array(F.array_distinct(F.flatten(F.array(
            F.coalesce(F.col("domain_spark_application_ids"), empty_string_array),
            F.coalesce(F.col("component_spark_application_ids"), empty_string_array),
            F.coalesce(F.col("action_spark_application_ids"), empty_string_array),
        )))),
    )
    .withColumn(
        "has_external_target",
        F.array_contains(F.col("target_scopes"), "Fabric 外部(確認済み)")
        | (
            F.size(F.array_intersect(
            F.coalesce(F.col("domain_destination_kinds"), empty_string_array),
                F.array(*[F.lit(x) for x in EXTERNAL_DESTINATION_KINDS]),
            )) > 0
        ),
    )
    .withColumn("has_outbound_action", F.size(F.coalesce(F.col("outbound_action_evidence"), empty_string_array)) > 0)
    .withColumn("has_component", F.size(F.coalesce(F.col("components_observed"), empty_string_array)) > 0)
    .withColumn(
        "has_unresolved_destination",
        F.array_contains(F.coalesce(F.col("domain_destination_kinds"), empty_string_array), "ドメイン確認対象"),
    )
    .withColumn(
        "has_runtime_service_client",
        F.size(F.coalesce(F.col("runtime_service_clients"), empty_string_array)) > 0,
    )
    .withColumn(
        "has_connector_or_driver",
        (
            F.size(F.coalesce(F.col("runtime_connectors_or_drivers"), empty_string_array)) > 0
        ) | (
            F.size(F.coalesce(F.col("connector_or_driver_dependencies"), empty_string_array)) > 0
        ),
    )
    .withColumn(
        "requires_destination_review",
        F.col("has_unresolved_destination")
        & F.col("has_runtime_service_client")
        & F.col("has_connector_or_driver"),
    )
)

detected = (
    base
    .withColumn("detect_target_date", F.to_date(F.col("submitted_datetime")))
    .withColumn(
        "detection_status",
        F.when(F.col("has_external_target") & F.col("has_outbound_action"), F.lit("Fabric 外部送信の根拠あり"))
         .when(F.col("has_external_target"), F.lit("Fabric 外部の接続先を検出"))
         .when(F.col("requires_destination_review"), F.lit("要確認"))
         .when(F.array_contains(F.col("target_scopes"), "Fabric 内部の可能性(要確認)"), F.lit("Fabric 内部の可能性(要確認)"))
         .when(F.col("has_component"), F.lit("外部送信の根拠なし"))
         .otherwise(F.lit("要確認")),
    )
    .withColumn(
        "detection_note",
        F.when(F.col("has_external_target") & F.col("has_outbound_action"), F.lit("同一実行内で Fabric 外部の接続先と具体的な書込み操作を検出した。保持した根拠で対象と操作の対応を確認する。"))
         .when(F.col("has_external_target"), F.lit("Fabric 外部の接続先を検出したが、具体的な書込み操作は検出されなかった。"))
         .when(F.col("requires_destination_review"), F.lit("接続先ドメイン、サービスクライアント、およびコネクター/ドライバーを検出したが、宛先種別を特定できないため要確認。"))
         .when(F.col("has_component"), F.lit("クライアント、コネクター、ドライバー、JAR または SDK を検出したが、これだけではデータ送信の証拠にならない。"))
         .otherwise(F.lit("保持した根拠を確認する。")),
    )
)

# 検知ロジック完了後に、結果に残った domain/client/connector の情報だけで確認種別を作る。
# 宛先を特定できた場合は確認ラベルより優先して destination 名を出力する。
known_destination_kinds = F.array_intersect(
    F.coalesce(F.col("domain_destination_kinds"), empty_string_array),
    F.array(*[F.lit(x) for x in EXTERNAL_DESTINATION_KINDS]),
)
confirmation_kinds = F.concat(
    F.when(
        F.size(F.coalesce(F.col("non_fabric_targets"), empty_string_array)) > 0,
        F.array(F.lit("利用したドメイン確認必要")),
    ).otherwise(empty_string_array),
    F.when(
        F.size(F.coalesce(F.col("runtime_service_clients"), empty_string_array)) > 0,
        F.array(F.lit("利用したクライアント確認必要")),
    ).otherwise(empty_string_array),
    F.when(
        F.size(F.coalesce(F.col("connector_or_driver_dependencies"), empty_string_array)) > 0,
        F.array(F.lit("利用したコネクター・ドライバー確認必要")),
    ).otherwise(empty_string_array),
)

result = (
    detected
    .withColumn(
        "destination_kinds",
        F.when(F.size(known_destination_kinds) > 0, known_destination_kinds)
         .when(F.size(confirmation_kinds) > 0, confirmation_kinds)
         .otherwise(F.lit(None).cast("array<string>")),
    )
    .where(F.col("destination_kinds").isNotNull())
    .select(
        *GROUP_KEYS, "detect_target_date", "spark_application_ids", "detection_status", "detection_note",
        "non_fabric_targets", "target_scopes", "destination_kinds", "non_fabric_domain_evidence",
        "runtime_service_clients", "connector_or_driver_dependencies",
        "outbound_action_kinds", "outbound_action_evidence",
    )
)

(result.write.format("delta").mode("overwrite").option("overwriteSchema", "true").saveAsTable(TARGET_TABLE))

display(result.orderBy("submitted_datetime", "workspace_id", "notebook_id", "livy_id"))

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_result_before_exclution

Để 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