0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

NHTSA FARS 交通事故データを Databricks (Free Edition) で分析する

0
Posted at

はじめに

過去の記事では、交通事故統計情報のオープンデータ (警察庁) を使用して日本の交通事故事情を分析しました。

今回は、海外の交通事故事情を分析してみます。米国運輸省 道路交通安全局(NHTSA)が非常にリッチな事故データをオープンデータとして公開しているので、個人的に気になっているデータプラットフォーム Databricks の Free Edition(無料枠)を利用して、米国交通死亡事故データの取り込みから加工、そして AI 機能を使った分析までを一気通貫で試してみたいと思います。

想定読者

  • モビリティ・交通安全領域のデータサイエンスに関心がある方
  • オープンデータを活用した実践的なデータ分析パイプラインに興味がある方
  • Databricks(特に最新のFree EditionやGenie Agents)を触ってみたい方

NHTSA FARS とは

死亡事故分析報告システム (Fatality Analysis Reporting System, FARS) は、米国運輸省 道路交通安全局 (NHTSA) が運用するシステムです。1975年から運用されており、全米の公道で発生した交通死亡事故の全数データが収集・公開されています。

なお、NHTSAは目的の異なる複数のデータシステムを併用しており、FARS以外にも以下のようなデータが存在します。

  • CISS (Crash Investigation Sampling System):
    • 抽出した事故について、車両損傷や傷害などを専門調査員が詳しく調査する詳細なサンプルデータ
  • CIREN (Crash Injury Research and Engineering Network):
    • 医療・工学の観点から、実際の事故における傷害発生メカニズムを詳細に研究するためのデータ
  • SCI (Special Crash Investigations):
    • 新技術や特殊な事故形態など、NHTSAが特に注目する事故を対象にした個別の詳細事故調査データ
  • FARS (Fatality Analysis Reporting System): ⇦ 本記事の対象
    • 全米の「交通死亡事故」の全数データ
  • CRSS (Crash Report Sampling System):
    • 警察報告された物損・負傷・死亡事故などから抽出した全国代表サンプル

また、CDAN (Crash Data Acquisition Network) は FARS / CRSS / CISS / SCI / CIREN など、NHTSAの複数の事故データ収集システムを支える共通のデータ収集・分析基盤です。

image.png
引用: 国土交通省の米国における交通事故調査に関する実態調査

これらのデータは、自動運転技術の安全性評価(Applied Intuition社のブログ等で言及)や、日本の国土交通省のレポートなどでも引用される、モビリティ業界における極めて重要な一次情報です。National Center for Statistics and Analysis (NCSA) Motor Vehicle Traffic Crash Data Resource Pageでは、これらを分析したレポートを見ることができます。

Databricks Free Edition とは

Databricks が提供する無料の学習・検証用環境です。大幅なアップデートが行われており、「無料版」と言いつつ、現在はできることが非常に多い強力な環境です。

Free Editionでは、NotebookやSQLによる分析だけでなく、モダンなデータ基盤を構築するための機能もかなり広く試せます。

  • Unity Catalog によるデータ・AI資産の管理
  • Serverless Compute によるNotebook / SQLの実行
  • Lakeflow によるデータ取り込み・変換パイプライン
  • Databricks AI/BI による可視化・分析
  • Genie Agents による自然言語でのデータ分析
  • Genie Code によるコーディング支援

今回は、これらの機能を利用して分析基盤を構築します。

※ 機能や利用上限はアップデートされる可能性があるため、最新情報は公式ドキュメントを確認してください。
https://www.databricks.com/jp/learn/free-edition
https://www.databricks.com/jp/blog/whats-coming-next-free-edition
https://www.databricks.com/jp/try-databricks

分析パイプラインの構築

モダンなデータエンジニアリングの基本である Ingestion (取り込み)Transformation (加工)Serving (提供) の 3 ステップで FARS データを処理していきます。

image.png
引用: Fundamentals of Data Engineering

Step 1: Ingestion

FARSには Crash Viewer API も用意されていますが、今回は一括でデータを取得してデータ基盤に取り込むためAPIは使わず、ZIPファイルのバルクダウンロードを利用します。

まずはデータの保存先となる Volume(Unity Catalog管理下のファイルストレージ)を作成します。

今回は workspace.default.volume です。
CatalogMyorganization/workspace/default へ進み、右上の CreateVolume から作成しました。

image.png

次に、FARSのサイトからデータを一括ダウンロードして解凍する Notebook を作成します。
+ NewNotebook で新規作成し、以下の Python コードを実行します。

00_FARS_DL_2024.ipynb
import os
import requests
from zipfile import ZipFile
from io import BytesIO

target_dir = "/Volumes/workspace/default/volume/FARS/"
year = "2024"
base_url = f"https://static.nhtsa.gov/nhtsa/downloads/FARS/{year}/National/"
filenames = [
    f"FARS{year}NationalCSV.zip", 
    f"FARS{year}NationalAuxiliaryCSV.zip", 
    f"FARS{year}NationalSAS.zip",
    f"FARS{year}NationalAuxiliarySAS.zip"
]
save_path = f"{target_dir}{year}/National/"

os.makedirs(save_path, exist_ok=True)

# リストの中身を一つずつダウンロード&解凍する
for fname in filenames:
    file_url = base_url + fname
    print(f"Downloading: {file_url}")    
    response = requests.get(file_url)
    
    if response.status_code == 200:
        try:
            with ZipFile(BytesIO(response.content)) as z:
                z.extractall(save_path)
                print(f"✅ Extracted: {fname} to {save_path}")
        except Exception as e:
            print(f"❌ Failed to extract {fname}: {e}")
    else:
        print(f"❌ Failed to download {fname}: Status Code {response.status_code}")

実行が完了すると、Volume内に以下のような構造でデータが取得できているはずです。

image.png

volume内のディレクトリ構造
/Volumes/workspace/default/volume/FARS/2024/National/
 ├── ACC_AUX.CSV                 <-- 補助データ
 ├── PER_AUX.CSV
 ├── VEH_AUX.CSV
 ├── (SASファイルたち...)
 └── FARS2024NationalCSV/        <-- メインのCSVデータはこのフォルダの中に!
      ├── accident.csv
      ├── person.csv
      ├── vehicle.csv
      ├── parkwork.csv
      └── ...その他のCSV

Step 2: Transformation

続いて ETL Pipeline (Lakeflow Pipelines) を作ります。

メダリオンアーキテクチャ(Bronze: 生データ → Silver: クレンジング → Gold: 分析用マート)を採用します。FARSでは事故・車両・人物でデータの粒度が異なるため、粒度を維持したまま加工します。また、キーの一意性や参照整合性、FARSの公式集計値との突合も行い、加工結果の妥当性を確認します。

+ NewETL pipeline で新規パイプラインを作成し、my_transformation.py を以下で書き換えました。

my_transformation.py
my_transformation.py
from pyspark import pipelines as dp
import pyspark.sql.functions as F


# ==========================================
# Configuration
# ==========================================

FARS_YEAR = 2024

BASE_PATH = (
    f"/Volumes/workspace/default/volume/"
    f"FARS/{FARS_YEAR}/National/**"
)


# ==========================================
# Helpers
# ==========================================

def prefixed_for_join(df, prefix, keys, match_col):
    """
    JOINキーはそのまま保持し、
    その他の列にはprefixを付与する。

    match_colは参照整合性確認用。
    """
    return df.select(
        *[F.col(k) for k in keys],
        F.lit(True).alias(match_col),
        *[
            F.col(c).alias(f"{prefix}{c}")
            for c in df.columns
            if c not in keys
        ]
    )


def pk_counts(df, dataset_name, keys):
    """
    主キー候補ごとの件数を返す。

    num_entries = 1
      → 一意

    num_entries > 1
      → 重複キー
    """
    return (
        df
        .groupBy(*keys)
        .count()
        .select(
            F.lit(dataset_name).alias("dataset"),
            F.concat_ws(
                "|",
                *[
                    F.col(k).cast("string")
                    for k in keys
                ]
            ).alias("key_value"),
            F.col("count").alias("num_entries")
        )
    )


# ==========================================
# 🥉 Bronze層
# Rawデータの取り込み
# ==========================================

# CSVはBronzeでは原則Stringとして取り込み、
# データ型・業務上の意味付けはSilverで行う。


# ------------------------------------------
# Accident
# ------------------------------------------

@dp.table(
    name="fars_accident_bronze",
    comment="FARS 2024 accident Rawデータ"
)
def fars_accident_bronze():

    return (
        spark.readStream
        .format("cloudFiles")
        .option("cloudFiles.format", "csv")
        .option("header", "true")
        .option(
            "cloudFiles.schemaEvolutionMode",
            "addNewColumns"
        )
        .option(
            "cloudFiles.useStrictGlobber",
            "true"
        )
        .option(
            "pathGlobFilter",
            "*[aA][cC][cC][iI][dD][eE][nN][tT].[cC][sS][vV]"
        )
        .load(BASE_PATH)

        .withColumn(
            "data_year",
            F.lit(FARS_YEAR)
        )
        .withColumn(
            "source_file_path",
            F.col("_metadata.file_path")
        )
        .withColumn(
            "ingestion_timestamp",
            F.current_timestamp()
        )
    )


# ------------------------------------------
# Vehicle
# Motor Vehicle In-Transport
# ------------------------------------------

@dp.table(
    name="fars_vehicle_bronze",
    comment="FARS 2024 vehicle Rawデータ"
)
def fars_vehicle_bronze():

    return (
        spark.readStream
        .format("cloudFiles")
        .option("cloudFiles.format", "csv")
        .option("header", "true")
        .option(
            "cloudFiles.schemaEvolutionMode",
            "addNewColumns"
        )
        .option(
            "cloudFiles.useStrictGlobber",
            "true"
        )
        .option(
            "pathGlobFilter",
            "*[vV][eE][hH][iI][cC][lL][eE].[cC][sS][vV]"
        )
        .load(BASE_PATH)

        .withColumn(
            "data_year",
            F.lit(FARS_YEAR)
        )
        .withColumn(
            "source_file_path",
            F.col("_metadata.file_path")
        )
        .withColumn(
            "ingestion_timestamp",
            F.current_timestamp()
        )
    )


# ------------------------------------------
# Parkwork
# Motor Vehicle Not In-Transport
# ------------------------------------------

@dp.table(
    name="fars_parkwork_bronze",
    comment="FARS 2024 parkwork Rawデータ"
)
def fars_parkwork_bronze():

    return (
        spark.readStream
        .format("cloudFiles")
        .option("cloudFiles.format", "csv")
        .option("header", "true")
        .option(
            "cloudFiles.schemaEvolutionMode",
            "addNewColumns"
        )
        .option(
            "cloudFiles.useStrictGlobber",
            "true"
        )
        .option(
            "pathGlobFilter",
            "*[pP][aA][rR][kK][wW][oO][rR][kK].[cC][sS][vV]"
        )
        .load(BASE_PATH)

        .withColumn(
            "data_year",
            F.lit(FARS_YEAR)
        )
        .withColumn(
            "source_file_path",
            F.col("_metadata.file_path")
        )
        .withColumn(
            "ingestion_timestamp",
            F.current_timestamp()
        )
    )


# ------------------------------------------
# Person
# ------------------------------------------

@dp.table(
    name="fars_person_bronze",
    comment="FARS 2024 person Rawデータ"
)
def fars_person_bronze():

    return (
        spark.readStream
        .format("cloudFiles")
        .option("cloudFiles.format", "csv")
        .option("header", "true")
        .option(
            "cloudFiles.schemaEvolutionMode",
            "addNewColumns"
        )
        .option(
            "cloudFiles.useStrictGlobber",
            "true"
        )
        .option(
            "pathGlobFilter",
            "*[pP][eE][rR][sS][oO][nN].[cC][sS][vV]"
        )
        .load(BASE_PATH)

        .withColumn(
            "data_year",
            F.lit(FARS_YEAR)
        )
        .withColumn(
            "source_file_path",
            F.col("_metadata.file_path")
        )
        .withColumn(
            "ingestion_timestamp",
            F.current_timestamp()
        )
    )


# ==========================================
# 🥈 Silver層
# 型・キーの正規化
# ==========================================


# ------------------------------------------
# Accident
#
# Grain:
# 1行 = 1事故
#
# PK:
# data_year + ST_CASE
# ------------------------------------------

@dp.table(
    name="fars_accident_silver",
    comment="型・キーを正規化した事故データ(1行 = 1事故)"
)
@dp.expect_or_fail(
    "valid_accident_key",
    """
    ST_CASE RLIKE '^[0-9]{6}$'
    AND data_year = 2024
    """
)
def fars_accident_silver():

    return (
        spark.readStream
        .table("fars_accident_bronze")

        # FARS Case Number
        .withColumn(
            "ST_CASE",
            F.lpad(
                F.trim(F.col("ST_CASE")),
                6,
                "0"
            )
        )

        # 分析用の型付き列
        .withColumn(
            "YEAR_VALUE",
            F.col("YEAR").cast("int")
        )
        .withColumn(
            "VE_TOTAL_VALUE",
            F.col("VE_TOTAL").cast("int")
        )
        .withColumn(
            "FATALS_VALUE",
            F.col("FATALS").cast("int")
        )
        .withColumn(
            "WEATHER_CODE",
            F.col("WEATHER").cast("int")
        )
        .withColumn(
            "LGT_COND_CODE",
            F.col("LGT_COND").cast("int")
        )
    )


# ------------------------------------------
# Vehicle
#
# Motor Vehicle In-Transport
#
# Grain:
# 1行 = 1車両
#
# PK:
# data_year + ST_CASE + VEH_NO
# ------------------------------------------

@dp.table(
    name="fars_vehicle_silver",
    comment="型・キーを正規化した走行中車両データ(1行 = 1車両)"
)
@dp.expect_or_fail(
    "valid_vehicle_key",
    """
    ST_CASE RLIKE '^[0-9]{6}$'
    AND data_year = 2024
    AND VEH_NO > 0
    """
)
def fars_vehicle_silver():

    return (
        spark.readStream
        .table("fars_vehicle_bronze")

        .withColumn(
            "ST_CASE",
            F.lpad(
                F.trim(F.col("ST_CASE")),
                6,
                "0"
            )
        )
        .withColumn(
            "VEH_NO",
            F.col("VEH_NO").cast("int")
        )

        .withColumn(
            "HIT_RUN_CODE",
            F.col("HIT_RUN").cast("int")
        )
    )


# ------------------------------------------
# Parkwork
#
# Motor Vehicle Not In-Transport
#
# Grain:
# 1行 = 1車両
#
# PK:
# data_year + ST_CASE + VEH_NO
# ------------------------------------------

@dp.table(
    name="fars_parkwork_silver",
    comment="型・キーを正規化した非走行中・作業中車両データ(1行 = 1車両)"
)
@dp.expect_or_fail(
    "valid_parkwork_key",
    """
    ST_CASE RLIKE '^[0-9]{6}$'
    AND data_year = 2024
    AND VEH_NO > 0
    """
)
def fars_parkwork_silver():

    return (
        spark.readStream
        .table("fars_parkwork_bronze")

        .withColumn(
            "ST_CASE",
            F.lpad(
                F.trim(F.col("ST_CASE")),
                6,
                "0"
            )
        )
        .withColumn(
            "VEH_NO",
            F.col("VEH_NO").cast("int")
        )
    )


# ------------------------------------------
# Person
#
# Grain:
# 1行 = 1人
#
# PK:
# data_year + ST_CASE + VEH_NO + PER_NO
# ------------------------------------------

@dp.table(
    name="fars_person_silver",
    comment="型・キーを正規化した人物データ(1行 = 1人)"
)
@dp.expect_or_fail(
    "valid_person_key",
    """
    ST_CASE RLIKE '^[0-9]{6}$'
    AND data_year = 2024
    AND VEH_NO >= 0
    AND PER_NO > 0
    """
)
@dp.expect_or_fail(
    "valid_person_vehicle_semantics",
    """
    (
        PER_TYP_CODE IN (1, 2, 3, 9)
        AND VEH_NO > 0
    )
    OR
    (
        PER_TYP_CODE IN (
            4, 5, 6, 7, 8,
            10, 11, 12, 13, 19
        )
        AND VEH_NO = 0
    )
    """
)
def fars_person_silver():

    return (
        spark.readStream
        .table("fars_person_bronze")

        .withColumn(
            "ST_CASE",
            F.lpad(
                F.trim(F.col("ST_CASE")),
                6,
                "0"
            )
        )
        .withColumn(
            "VEH_NO",
            F.col("VEH_NO").cast("int")
        )
        .withColumn(
            "PER_NO",
            F.col("PER_NO").cast("int")
        )

        # Person Type
        .withColumn(
            "PER_TYP_CODE",
            F.col("PER_TYP").cast("int")
        )

        # Non-motorist striking vehicle
        #
        # Generic Person GoldではJOINせず、
        # Person属性として保持する。
        .withColumn(
            "STR_VEH_CODE",
            F.col("STR_VEH").cast("int")
        )

        .withColumn(
            "AGE_VALUE",
            F.col("AGE").cast("int")
        )
        .withColumn(
            "INJ_SEV_CODE",
            F.col("INJ_SEV").cast("int")
        )
    )


# ==========================================
# 🔎 Data Quality 1
# Silver PK uniqueness
# ==========================================

@dp.materialized_view(
    name="fars_pk_check",
    comment="Silverの主キー一意性チェック"
)
@dp.expect_or_fail(
    "unique_primary_key",
    "num_entries = 1"
)
def fars_pk_check():

    accident = pk_counts(
        spark.read.table(
            "fars_accident_silver"
        ),
        "accident_silver",
        [
            "data_year",
            "ST_CASE"
        ]
    )

    vehicle = pk_counts(
        spark.read.table(
            "fars_vehicle_silver"
        ),
        "vehicle_silver",
        [
            "data_year",
            "ST_CASE",
            "VEH_NO"
        ]
    )

    parkwork = pk_counts(
        spark.read.table(
            "fars_parkwork_silver"
        ),
        "parkwork_silver",
        [
            "data_year",
            "ST_CASE",
            "VEH_NO"
        ]
    )

    person = pk_counts(
        spark.read.table(
            "fars_person_silver"
        ),
        "person_silver",
        [
            "data_year",
            "ST_CASE",
            "VEH_NO",
            "PER_NO"
        ]
    )

    return (
        accident
        .unionByName(vehicle)
        .unionByName(parkwork)
        .unionByName(person)
    )


# ==========================================
# 🥇 Gold層
# Grain-specific Data Marts
# ==========================================


# ------------------------------------------
# Gold 1: Accident
#
# Grain:
# 1行 = 1事故
# ------------------------------------------

@dp.materialized_view(
    name="fars_accident_summary_gold",
    comment="【事故粒度】1行=1事故の分析用データマート"
)
def fars_accident_summary_gold():

    acc_df = spark.read.table(
        "fars_accident_silver"
    )

    per_df = spark.read.table(
        "fars_person_silver"
    )

    # Personから事故単位の人物数を集計
    per_agg = (
        per_df
        .groupBy(
            "data_year",
            "ST_CASE"
        )
        .agg(
            F.count("*")
            .alias("total_persons_involved")
        )
    )

    return (
        acc_df

        # 車両数・死亡者数は
        # Accident側の公式集計値を利用
        .withColumn(
            "total_vehicles_involved",
            F.col("VE_TOTAL_VALUE")
        )
        .withColumn(
            "total_fatalities",
            F.col("FATALS_VALUE")
        )

        .join(
            per_agg,
            on=[
                "data_year",
                "ST_CASE"
            ],
            how="left"
        )
    )


# ------------------------------------------
# Gold 2: Vehicle
#
# Grain:
# 1行 = 1 Motor Vehicle In-Transport
# ------------------------------------------

@dp.materialized_view(
    name="fars_vehicle_wide_gold",
    comment="【車両粒度】走行中車両を主語に事故情報を付与した分析用データマート"
)
@dp.expect_or_fail(
    "vehicle_has_accident",
    "vehicle_accident_reference_ok = true"
)
def fars_vehicle_wide_gold():

    veh_df = spark.read.table(
        "fars_vehicle_silver"
    )

    acc_df = spark.read.table(
        "fars_accident_silver"
    )

    join_keys = [
        "data_year",
        "ST_CASE"
    ]

    accident = prefixed_for_join(
        acc_df,
        prefix="acc_",
        keys=join_keys,
        match_col="_acc_match"
    )

    return (
        veh_df
        .join(
            accident,
            on=join_keys,
            how="left"
        )

        .withColumn(
            "vehicle_accident_reference_ok",
            F.coalesce(
                F.col("_acc_match"),
                F.lit(False)
            )
        )

        .drop("_acc_match")
    )


# ------------------------------------------
# Gold 3: Person
#
# Grain:
# 1行 = 1人
# ------------------------------------------

@dp.materialized_view(
    name="fars_person_wide_gold",
    comment="【人物粒度】人物を主語に事故・所属車両情報を付与した分析用データマート"
)
@dp.expect_or_fail(
    "person_has_accident",
    "person_accident_reference_ok = true"
)
@dp.expect_or_fail(
    "valid_person_vehicle_reference",
    "person_vehicle_reference_ok = true"
)
def fars_person_wide_gold():

    per_df = spark.read.table(
        "fars_person_silver"
    )

    veh_df = spark.read.table(
        "fars_vehicle_silver"
    )

    park_df = spark.read.table(
        "fars_parkwork_silver"
    )

    acc_df = spark.read.table(
        "fars_accident_silver"
    )


    # ======================================
    # Person側に参照キーを定義
    # ======================================

    person_base = (
        per_df

        # PER_TYP 1, 2, 9
        # Motor Vehicle In-Transport occupant
        .withColumn(
            "occupant_vehicle_no",
            F.when(
                F.col("PER_TYP_CODE")
                .isin(1, 2, 9),
                F.col("VEH_NO")
            )
        )

        # PER_TYP 3
        # Motor Vehicle Not In-Transport occupant
        .withColumn(
            "parkwork_vehicle_no",
            F.when(
                F.col("PER_TYP_CODE") == 3,
                F.col("VEH_NO")
            )
        )
    )


    # ======================================
    # ① Motor Vehicle In-Transport
    # ======================================

    occupant_vehicle = (
        veh_df
        .select(
            "data_year",
            "ST_CASE",

            F.col("VEH_NO").alias(
                "occupant_vehicle_no"
            ),

            F.lit(True).alias(
                "_occupant_vehicle_match"
            ),

            *[
                F.col(c).alias(
                    f"occ_veh_{c}"
                )
                for c in veh_df.columns
                if c not in [
                    "data_year",
                    "ST_CASE",
                    "VEH_NO"
                ]
            ]
        )
    )

    person_base = (
        person_base
        .join(
            occupant_vehicle,
            on=[
                "data_year",
                "ST_CASE",
                "occupant_vehicle_no"
            ],
            how="left"
        )
    )


    # ======================================
    # ② Motor Vehicle Not In-Transport
    # ======================================

    parkwork_vehicle = (
        park_df
        .select(
            "data_year",
            "ST_CASE",

            F.col("VEH_NO").alias(
                "parkwork_vehicle_no"
            ),

            F.lit(True).alias(
                "_parkwork_vehicle_match"
            ),

            *[
                F.col(c).alias(
                    f"park_veh_{c}"
                )
                for c in park_df.columns
                if c not in [
                    "data_year",
                    "ST_CASE",
                    "VEH_NO"
                ]
            ]
        )
    )

    person_base = (
        person_base
        .join(
            parkwork_vehicle,
            on=[
                "data_year",
                "ST_CASE",
                "parkwork_vehicle_no"
            ],
            how="left"
        )
    )


    # ======================================
    # ③ Accident
    # ======================================

    acc_keys = [
        "data_year",
        "ST_CASE"
    ]

    accident = prefixed_for_join(
        acc_df,
        prefix="acc_",
        keys=acc_keys,
        match_col="_acc_match"
    )

    person_base = (
        person_base
        .join(
            accident,
            on=acc_keys,
            how="left"
        )
    )


    # ======================================
    # ④ Referential Integrity
    # ======================================

    return (
        person_base

        # Vehicle / Parkworkへの参照整合性
        .withColumn(
            "person_vehicle_reference_ok",

            F.when(
                F.col("PER_TYP_CODE")
                .isin(1, 2, 9),

                F.coalesce(
                    F.col(
                        "_occupant_vehicle_match"
                    ),
                    F.lit(False)
                )
            )

            .when(
                F.col("PER_TYP_CODE") == 3,

                F.coalesce(
                    F.col(
                        "_parkwork_vehicle_match"
                    ),
                    F.lit(False)
                )
            )

            # Non-motoristはGeneric Person Goldでは
            # VehicleへJOINしない
            .otherwise(
                F.lit(True)
            )
        )

        # Accidentへの参照整合性
        .withColumn(
            "person_accident_reference_ok",

            F.coalesce(
                F.col("_acc_match"),
                F.lit(False)
            )
        )

        .drop(
            "_occupant_vehicle_match",
            "_parkwork_vehicle_match",
            "_acc_match"
        )
    )


# ==========================================
# 🔎 Data Quality 2
# Gold PK uniqueness
# ==========================================

@dp.materialized_view(
    name="fars_gold_pk_check",
    comment="Goldデータマートの粒度・主キー一意性チェック"
)
@dp.expect_or_fail(
    "unique_gold_key",
    "num_entries = 1"
)
def fars_gold_pk_check():

    accident = pk_counts(
        spark.read.table(
            "fars_accident_summary_gold"
        ),
        "accident_gold",
        [
            "data_year",
            "ST_CASE"
        ]
    )

    vehicle = pk_counts(
        spark.read.table(
            "fars_vehicle_wide_gold"
        ),
        "vehicle_gold",
        [
            "data_year",
            "ST_CASE",
            "VEH_NO"
        ]
    )

    person = pk_counts(
        spark.read.table(
            "fars_person_wide_gold"
        ),
        "person_gold",
        [
            "data_year",
            "ST_CASE",
            "VEH_NO",
            "PER_NO"
        ]
    )

    return (
        accident
        .unionByName(vehicle)
        .unionByName(person)
    )


# ==========================================
# 🔎 Data Quality 3
# Accident-level Reconciliation
# ==========================================

@dp.materialized_view(
    name="fars_reconciliation_check",
    comment="事故単位でFARS公式集計値とSilver明細を突合"
)
@dp.expect_or_fail(
    "fatalities_reconcile",
    "fatality_reconcile_ok = true"
)
@dp.expect_or_fail(
    "vehicles_reconcile",
    "vehicle_reconcile_ok = true"
)
def fars_reconciliation_check():

    # --------------------------------------
    # Accident側の公式値
    # --------------------------------------

    accident = (
        spark.read
        .table("fars_accident_silver")
        .select(
            "data_year",
            "ST_CASE",
            "FATALS_VALUE",
            "VE_TOTAL_VALUE"
        )
    )


    # --------------------------------------
    # Person
    #
    # 事故単位の死亡者数
    # --------------------------------------

    person = (
        spark.read
        .table("fars_person_silver")

        .groupBy(
            "data_year",
            "ST_CASE"
        )

        .agg(
            F.sum(
                F.when(
                    F.col("INJ_SEV_CODE") == 4,
                    1
                ).otherwise(0)
            ).alias(
                "person_fatalities"
            )
        )
    )


    # --------------------------------------
    # Vehicle
    #
    # 事故単位の
    # Motor Vehicle In-Transport数
    # --------------------------------------

    vehicle = (
        spark.read
        .table("fars_vehicle_silver")

        .groupBy(
            "data_year",
            "ST_CASE"
        )

        .agg(
            F.count("*").alias(
                "vehicle_rows"
            )
        )
    )


    # --------------------------------------
    # Parkwork
    #
    # 事故単位の
    # Motor Vehicle Not In-Transport数
    # --------------------------------------

    parkwork = (
        spark.read
        .table("fars_parkwork_silver")

        .groupBy(
            "data_year",
            "ST_CASE"
        )

        .agg(
            F.count("*").alias(
                "parkwork_rows"
            )
        )
    )


    # --------------------------------------
    # Accidentを基準に事故単位で突合
    # --------------------------------------

    return (
        accident

        .join(
            person,
            on=[
                "data_year",
                "ST_CASE"
            ],
            how="left"
        )

        .join(
            vehicle,
            on=[
                "data_year",
                "ST_CASE"
            ],
            how="left"
        )

        .join(
            parkwork,
            on=[
                "data_year",
                "ST_CASE"
            ],
            how="left"
        )

        # 明細レコードが存在しない場合は0件
        .fillna(
            0,
            subset=[
                "person_fatalities",
                "vehicle_rows",
                "parkwork_rows"
            ]
        )

        # ----------------------------------
        # FATALS reconciliation
        # ----------------------------------

        .withColumn(
            "fatality_reconcile_ok",

            F.col("FATALS_VALUE")
            == F.col("person_fatalities")
        )

        # ----------------------------------
        # VE_TOTAL reconciliation
        # ----------------------------------

        .withColumn(
            "derived_vehicle_total",

            F.col("vehicle_rows")
            + F.col("parkwork_rows")
        )

        .withColumn(
            "vehicle_reconcile_ok",

            F.col("VE_TOTAL_VALUE")
            == F.col("derived_vehicle_total")
        )
    )


# ==========================================
# 🔎 Data Quality 4
# Silver → Gold Grain Preservation
# ==========================================

@dp.materialized_view(
    name="fars_gold_row_count_check",
    comment="Gold生成時に粒度が崩れていないことを行数で検証"
)
@dp.expect_or_fail(
    "accident_rows_preserved",
    "accident_silver_rows = accident_gold_rows"
)
@dp.expect_or_fail(
    "vehicle_rows_preserved",
    "vehicle_silver_rows = vehicle_gold_rows"
)
@dp.expect_or_fail(
    "person_rows_preserved",
    "person_silver_rows = person_gold_rows"
)
def fars_gold_row_count_check():

    accident_silver = (
        spark.read
        .table("fars_accident_silver")
        .agg(
            F.count("*").alias(
                "accident_silver_rows"
            )
        )
    )

    accident_gold = (
        spark.read
        .table("fars_accident_summary_gold")
        .agg(
            F.count("*").alias(
                "accident_gold_rows"
            )
        )
    )


    vehicle_silver = (
        spark.read
        .table("fars_vehicle_silver")
        .agg(
            F.count("*").alias(
                "vehicle_silver_rows"
            )
        )
    )

    vehicle_gold = (
        spark.read
        .table("fars_vehicle_wide_gold")
        .agg(
            F.count("*").alias(
                "vehicle_gold_rows"
            )
        )
    )


    person_silver = (
        spark.read
        .table("fars_person_silver")
        .agg(
            F.count("*").alias(
                "person_silver_rows"
            )
        )
    )

    person_gold = (
        spark.read
        .table("fars_person_wide_gold")
        .agg(
            F.count("*").alias(
                "person_gold_rows"
            )
        )
    )


    # 各集計結果は1行なので
    # crossJoinして1行のvalidation datasetにする
    return (
        accident_silver
        .crossJoin(accident_gold)
        .crossJoin(vehicle_silver)
        .crossJoin(vehicle_gold)
        .crossJoin(person_silver)
        .crossJoin(person_gold)
    )

Run pipelineを押下すると、DAGが表示され、各テーブルが作成されます。

image.png

Gold層では、事故・車両・人物ごとに「1行が何を表すか」を明確にした3つのデータマートを作成しています。また、主キーの一意性やJOIN前後の件数、FARS公式値との整合性もPipeline内で検証しています。

※ 本記事では FARS 2024 National データを対象としています。他年度を利用する場合は、対象年度のデータ仕様を確認してください。

Step 3: Serving

いよいよ分析です。今回は Genie Agents で分析します。
※ 本記事執筆時点では、Free Editionで利用できました。

+ NewGenie Agent と進むと以下の画面になるので、前ステップで作成したGoldテーブルを指定します。

image.png

ここまでくれば、自然言語による分析ができると思います。
例えば「死亡事故について分析してください」と入力すると以下のような回答が返ってきます。

全体
image.png

時間的パターン
image.png

曜日別傾向
image.png

性別や飲酒状況など、主要な知見
image.png

Genie Agents の出力結果が必ずしも妥当とは限らないため、実際に利用する際は生成されたSQLや集計条件、元データの定義を確認してください。

おわりに

今回は NHTSA の FARS データを使用し、Databricks Free Edition 上でデータ取り込みから AI エージェントによる分析までを一気通貫で試してみました。オープンデータと無料のクラウド環境を組み合わせるだけで、サクッと最低限の分析基盤を構築できるのは、素晴らしいと感じました。

モビリティ領域のデータ分析に興味がある方は、ぜひ以下の関連リンク・データソースも覗いてみてください。さらなる分析のアイデアが見つかるはずです。

関連データソース・参考リンク
0
0
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?