はじめに
過去の記事では、交通事故統計情報のオープンデータ (警察庁) を使用して日本の交通事故事情を分析しました。
今回は、海外の交通事故事情を分析してみます。米国運輸省 道路交通安全局(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の複数の事故データ収集システムを支える共通のデータ収集・分析基盤です。
これらのデータは、自動運転技術の安全性評価(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 データを処理していきます。
Step 1: Ingestion
FARSには Crash Viewer API も用意されていますが、今回は一括でデータを取得してデータ基盤に取り込むためAPIは使わず、ZIPファイルのバルクダウンロードを利用します。
まずはデータの保存先となる Volume(Unity Catalog管理下のファイルストレージ)を作成します。
今回は workspace.default.volume です。
Catalog → Myorganization/workspace/default へ進み、右上の Create → Volume から作成しました。
次に、FARSのサイトからデータを一括ダウンロードして解凍する Notebook を作成します。
+ New → Notebook で新規作成し、以下の Python コードを実行します。
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内に以下のような構造でデータが取得できているはずです。
/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の公式集計値との突合も行い、加工結果の妥当性を確認します。
+ New → ETL pipeline で新規パイプラインを作成し、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が表示され、各テーブルが作成されます。
Gold層では、事故・車両・人物ごとに「1行が何を表すか」を明確にした3つのデータマートを作成しています。また、主キーの一意性やJOIN前後の件数、FARS公式値との整合性もPipeline内で検証しています。
※ 本記事では FARS 2024 National データを対象としています。他年度を利用する場合は、対象年度のデータ仕様を確認してください。
Step 3: Serving
いよいよ分析です。今回は Genie Agents で分析します。
※ 本記事執筆時点では、Free Editionで利用できました。
+ New → Genie Agent と進むと以下の画面になるので、前ステップで作成したGoldテーブルを指定します。
ここまでくれば、自然言語による分析ができると思います。
例えば「死亡事故について分析してください」と入力すると以下のような回答が返ってきます。
Genie Agents の出力結果が必ずしも妥当とは限らないため、実際に利用する際は生成されたSQLや集計条件、元データの定義を確認してください。
おわりに
今回は NHTSA の FARS データを使用し、Databricks Free Edition 上でデータ取り込みから AI エージェントによる分析までを一気通貫で試してみました。オープンデータと無料のクラウド環境を組み合わせるだけで、サクッと最低限の分析基盤を構築できるのは、素晴らしいと感じました。
モビリティ領域のデータ分析に興味がある方は、ぜひ以下の関連リンク・データソースも覗いてみてください。さらなる分析のアイデアが見つかるはずです。
関連データソース・参考リンク
-
Fatality Analysis Reporting System (FARS)
https://www.nhtsa.gov/research-data/fatality-analysis-reporting-system-fars -
NHTSA Crash Data Systems
https://www.nhtsa.gov/data/crash-data-systems -
National Center for Statistics and Analysis (NCSA)
https://www.nhtsa.gov/research-data/national-center-statistics-and-analysis-ncsa -
NCSA | Tools, Publications, and Data
https://cdan.dot.gov/ -
FARS 2024 National File Download
https://www.nhtsa.gov/file-downloads?p=nhtsa/downloads/FARS/2024/National/ -
CrashStats Publications
https://crashstats.nhtsa.dot.gov/#!/PublicationList/51 -
OECD/ITF Road Safety Country Profiles
https://www.itf-oecd.org/road-safety-country-profiles -
OECD/ITF Road Safety Dashboard
https://www.itf-oecd.org/road-safety-dashboard -
国土交通省 | 車両安全対策検討会
https://www.mlit.go.jp/jidosha/jidosha_tk7_000005.html -
国土交通省 | 米国における交通事故調査に関する実態調査
https://www.mlit.go.jp/jidosha/content/001843919.pdf









