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