Commit 3f14d770a2
Verified · cmc
Layout: unified · split
.github/workflows/daily-pull.yml added +151
| @@ -0,0 +1,151 @@ | ||
| 1 | # Sarpy and Council Bluffs serve a rolling 12-month window; records that age out | |
| 2 | # of those feeds exist nowhere else. This job is the only thing keeping them, so | |
| 3 | # it refuses to publish an archive smaller than the one it started with. | |
| 4 | name: daily pull | |
| 5 | ||
| 6 | on: | |
| 7 | schedule: | |
| 8 | # 06:00 America/Chicago in summer, 05:00 in winter. Offset from the hour | |
| 9 | # because GitHub drops on-the-hour scheduled runs under load. | |
| 10 | - cron: "17 11 * * *" | |
| 11 | workflow_dispatch: | |
| 12 | inputs: | |
| 13 | bootstrap: | |
| 14 | description: "Start a new archive instead of restoring the published one" | |
| 15 | type: boolean | |
| 16 | default: false | |
| 17 | full: | |
| 18 | description: "Pull each feed in full rather than the last 30 days" | |
| 19 | type: boolean | |
| 20 | default: false | |
| 21 | ||
| 22 | permissions: | |
| 23 | contents: write | |
| 24 | ||
| 25 | concurrency: | |
| 26 | group: archive | |
| 27 | cancel-in-progress: false | |
| 28 | ||
| 29 | env: | |
| 30 | TAG: archive | |
| 31 | DB: raw_data/metro.db | |
| 32 | GH_TOKEN: ${{ github.token }} | |
| 33 | ||
| 34 | jobs: | |
| 35 | pull: | |
| 36 | runs-on: ubuntu-latest | |
| 37 | timeout-minutes: 30 | |
| 38 | ||
| 39 | steps: | |
| 40 | - uses: actions/checkout@v4 | |
| 41 | ||
| 42 | - uses: actions/setup-python@v5 | |
| 43 | with: | |
| 44 | python-version: "3.13" | |
| 45 | cache: pip | |
| 46 | cache-dependency-path: requirements-ingest.txt | |
| 47 | ||
| 48 | - run: pip install -r requirements-ingest.txt | |
| 49 | ||
| 50 | - name: Restore archive | |
| 51 | run: | | |
| 52 | if gh release download "$TAG" --pattern metro.db.gz --dir .; then | |
| 53 | gunzip -c metro.db.gz > "$DB" | |
| 54 | rm metro.db.gz | |
| 55 | echo "restored $(du -h "$DB" | cut -f1)" | |
| 56 | elif [ "${{ inputs.bootstrap }}" = "true" ]; then | |
| 57 | echo "no published archive; starting a new one" | |
| 58 | else | |
| 59 | echo "::error::No metro.db.gz on release '$TAG'. Anything that has" \ | |
| 60 | "already aged out of the Sarpy and Council Bluffs feeds cannot" \ | |
| 61 | "be recovered. Re-run with bootstrap only if that is intended." | |
| 62 | exit 1 | |
| 63 | fi | |
| 64 | ||
| 65 | - name: Count rows before | |
| 66 | run: | | |
| 67 | : > before.txt | |
| 68 | if [ -f "$DB" ]; then | |
| 69 | sqlite3 -noheader -separator ' ' "$DB" \ | |
| 70 | "SELECT source, COUNT(*) FROM incidents GROUP BY source ORDER BY source" \ | |
| 71 | > before.txt | |
| 72 | fi | |
| 73 | cat before.txt | |
| 74 | ||
| 75 | - name: Pull feeds | |
| 76 | run: | | |
| 77 | if [ "${{ inputs.full }}" = "true" ] || [ "${{ inputs.bootstrap }}" = "true" ]; then | |
| 78 | python ingest.py --full | |
| 79 | else | |
| 80 | python ingest.py | |
| 81 | fi | |
| 82 | ||
| 83 | - name: Check nothing was lost | |
| 84 | run: | | |
| 85 | sqlite3 -noheader -separator ' ' "$DB" \ | |
| 86 | "SELECT source, COUNT(*) FROM incidents GROUP BY source ORDER BY source" \ | |
| 87 | > after.txt | |
| 88 | cat after.txt | |
| 89 | test -s after.txt || { echo "::error::archive is empty"; exit 1; } | |
| 90 | # Keyed on FILENAME, not NR == FNR: before.txt is empty on a bootstrap | |
| 91 | # run, and awk never resets FNR for a zero-length file. | |
| 92 | awk -v first=before.txt ' | |
| 93 | FILENAME == first { was[$1] = $2; next } | |
| 94 | { now[$1] = $2 } | |
| 95 | END { | |
| 96 | for (s in was) | |
| 97 | if (now[s] + 0 < was[s] + 0) { | |
| 98 | printf "::error::%s lost rows: %d -> %d\n", s, was[s], now[s] | |
| 99 | bad = 1 | |
| 100 | } | |
| 101 | exit bad | |
| 102 | }' before.txt after.txt | |
| 103 | ||
| 104 | - name: Publish archive | |
| 105 | run: | | |
| 106 | sqlite3 "$DB" "VACUUM;" | |
| 107 | gzip -c "$DB" > metro.db.gz | |
| 108 | gh release view "$TAG" >/dev/null 2>&1 \ | |
| 109 | || gh release create "$TAG" --title "Incident archive" --notes "building" | |
| 110 | gh release upload "$TAG" metro.db.gz --clobber | |
| 111 | { | |
| 112 | echo "SQLite archive of Omaha metro police incident feeds, rebuilt daily." | |
| 113 | echo "Sarpy County and Council Bluffs publish a rolling 12-month window," | |
| 114 | echo "so this holds records their own feeds no longer serve." | |
| 115 | echo | |
| 116 | echo "Updated $(date -u '+%Y-%m-%d %H:%M UTC'). Schema: schema.sql." | |
| 117 | echo | |
| 118 | echo '```' | |
| 119 | sqlite3 -header -column "$DB" \ | |
| 120 | "SELECT agency, COUNT(*) AS rows, SUM(is_stop) AS stops, | |
| 121 | MIN(occurred_at) AS earliest, MAX(occurred_at) AS latest | |
| 122 | FROM incidents GROUP BY agency ORDER BY rows DESC" | |
| 123 | echo '```' | |
| 124 | } > notes.md | |
| 125 | gh release edit "$TAG" --notes-file notes.md | |
| 126 | ||
| 127 | # Runs after the upload on purpose: a feed that stopped updating should | |
| 128 | # raise the alarm without also blocking the archive from being published. | |
| 129 | - name: Check the feeds are still moving | |
| 130 | run: | | |
| 131 | sqlite3 -noheader "$DB" \ | |
| 132 | "SELECT source || ' ' || MAX(occurred_at) FROM incidents | |
| 133 | WHERE source <> 'opd_csv' GROUP BY source | |
| 134 | HAVING MAX(occurred_at) < datetime('now', '-7 days')" > stale.txt | |
| 135 | if [ -s stale.txt ]; then | |
| 136 | while read -r line; do echo "::error::feed is stale: $line"; done < stale.txt | |
| 137 | exit 1 | |
| 138 | fi | |
| 139 | echo "all feeds current" | |
| 140 | ||
| 141 | - name: Summary | |
| 142 | if: always() | |
| 143 | run: | | |
| 144 | { | |
| 145 | echo "| source | before | after |" | |
| 146 | echo "|---|---|---|" | |
| 147 | awk -v first=before.txt ' | |
| 148 | FILENAME == first { was[$1] = $2; next } | |
| 149 | { printf "| %s | %s | %s |\n", $1, ($1 in was ? was[$1] : 0), $2 }' \ | |
| 150 | before.txt after.txt | |
| 151 | } >> "$GITHUB_STEP_SUMMARY" | |
.gitignore +3
| @@ -1 +1,4 @@ | ||
| 1 | 1 | .DS_Store |
| 2 | *.db | |
| 3 | .venv/ | |
| 4 | __pycache__/ | |
README.nfo +80 −14
| @@ -3,25 +3,91 @@ | ||
| 3 | 3 | └──────────────────────────────────────────────────────────────┘ |
| 4 | 4 | |
| 5 | 5 | WHAT |
| 6 | crime data from the omaha police department, for analysis and | |
| 7 | visualization. | |
| 6 | police activity across the omaha metro, pulled from the agencies' | |
| 7 | own feeds and joined against alpr camera locations from | |
| 8 | openstreetmap. | |
| 8 | 9 | |
| 9 | API | |
| 10 | explore the database over http with sqlite2rest: | |
| 10 | COVERAGE | |
| 11 | omaha pd dcgis arcgis view, 2022-01-01 onward, | |
| 12 | nibrs offence records, updated daily. | |
| 13 | no stop or disposition data. | |
| 14 | bellevue pd sarpy county publiccrimemap, cad calls | |
| 15 | papillion pd for service with stop type, disposition | |
| 16 | la vista pd and category. rolling 12-month window: | |
| 17 | sarpy county so records age out of the feed, so the local | |
| 18 | archive is the only long-term copy. gretna | |
| 19 | and springfield appear as fire only; their | |
| 20 | police departments do not report to it. | |
| 21 | council bluffs pd cbpd public cfs feed, refreshed every ten | |
| 22 | minutes, rolling 12-month window. stop | |
| 23 | type, disposition, priority, response time. | |
| 24 | street addresses withheld; points exact. | |
| 25 | ralston pd no machine-readable feed. absent. | |
| 11 | 26 | |
| 12 | pip3 install sqlite2rest | |
| 13 | sqlite2rest serve ./raw_data/ingress.db | |
| 27 | alpr cameras openstreetmap via overpass, the same data | |
| 28 | deflock renders. 168 nodes in the metro | |
| 29 | bbox. odbl, attribution required. | |
| 14 | 30 | |
| 15 | TODO | |
| 16 | - import script (done) | |
| 17 | - drop duplicate header rows (done) | |
| 18 | - explore and analyze the data | |
| 19 | - plotly dash visualizations | |
| 20 | - api to the database | |
| 31 | SETUP | |
| 32 | uv venv .venv | |
| 33 | uv pip install --python .venv/bin/python -r requirements.txt | |
| 34 | ||
| 35 | USE | |
| 36 | .venv/bin/python ingest.py --full # first run, backfill | |
| 37 | .venv/bin/python ingest.py # daily, last 30 days | |
| 38 | .venv/bin/python ingest.py cbpd sarpy # one source at a time | |
| 39 | .venv/bin/python ingest.py opd_csv # 2015-2023 csv archive | |
| 40 | .venv/bin/python app.py | |
| 41 | ||
| 42 | ARCHIVE | |
| 43 | .github/workflows/daily-pull.yml runs the pull at 11:17 utc and | |
| 44 | keeps the database as metro.db.gz on the "archive" release, so the | |
| 45 | archive does not depend on any one machine. each run restores that | |
| 46 | asset, pulls, refuses to publish if any source came back with fewer | |
| 47 | rows than it started with, then uploads and fails loudly if a feed | |
| 48 | has not moved in seven days. | |
| 49 | ||
| 50 | first run: trigger it manually with bootstrap enabled, which pulls | |
| 51 | every feed in full and creates the release. after that the restore | |
| 52 | step is mandatory -- a bootstrap over a live archive throws away | |
| 53 | whatever has already aged out of the sarpy and council bluffs | |
| 54 | feeds. | |
| 55 | ||
| 56 | github disables scheduled workflows after 60 days without repo | |
| 57 | activity, and emails first. that is the most likely way this stops | |
| 58 | quietly. | |
| 59 | ||
| 60 | to run the pull locally instead: | |
| 61 | ||
| 62 | 0 6 * * * cd /path/to/omaha-incidents && .venv/bin/python ingest.py | |
| 63 | ||
| 64 | NOTES | |
| 65 | all three arcgis services return utc epochs; their where-clause | |
| 66 | literals do not agree (opd and council bluffs utc, sarpy central). | |
| 67 | ingest.py stores occurred_at in local time. | |
| 68 | ||
| 69 | each feed has its own taxonomy and none of them are comparable, so | |
| 70 | ingest.py derives one cross-agency flag, is_stop, per source. stop | |
| 71 | outcomes compare citation and arrest rate by substring, which is | |
| 72 | all the two disposition vocabularies support: an agency that | |
| 73 | records warnings less thoroughly shows a higher citation rate for | |
| 74 | that reason alone. | |
| 75 | ||
| 76 | colour scheme follows prefers-color-scheme. plotly writes colours | |
| 77 | into the figure, so assets/theme.js reports the media query into a | |
| 78 | store and app.py builds each figure from it. restyling after the | |
| 79 | fact does not work: swapping a maplibre basemap at runtime leaves | |
| 80 | it rebuilding with no data layers. | |
| 81 | ||
| 82 | the camera-proximity panel compares stops against a non-stop | |
| 83 | baseline. cameras and stops both concentrate on arterials, so a | |
| 84 | gap between the curves is a starting point, not a finding. | |
| 85 | ||
| 86 | raw_data/ingress.db is the old 2015-2023 sqlite build. nothing | |
| 87 | reads it any more. | |
| 21 | 88 | |
| 22 | 89 | SCREENSHOTS |
| 23 | screenshots/dashboard_01.png | |
| 24 | screenshots/dashboard_02.png | |
| 90 | screenshots/*.png are from the previous 2015-2023 dashboard. | |
| 25 | 91 | |
| 26 | 92 | ┌──────────────────────────────────────────────────────────────┐ |
| 27 | 93 | │ krz.sh │ |
analysis.py added +146
| @@ -0,0 +1,146 @@ | ||
| 1 | """Queries behind the dashboard: agency activity and distance to the nearest ALPR.""" | |
| 2 | ||
| 3 | import sqlite3 | |
| 4 | from pathlib import Path | |
| 5 | ||
| 6 | import numpy as np | |
| 7 | import pandas as pd | |
| 8 | ||
| 9 | DB = Path(__file__).parent / "raw_data" / "metro.db" | |
| 10 | ||
| 11 | # Sarpy and Council Bluffs publish officer-initiated stops; ingest.py flags them | |
| 12 | # as is_stop. OPD publishes NIBRS offences only and contributes no stops. | |
| 13 | POLICE_AGENCIES = ("Omaha PD", "Council Bluffs PD", "Bellevue PD", "Papillion PD", | |
| 14 | "La Vista PD", "Sarpy County SO") | |
| 15 | ||
| 16 | EARTH_M = 6371000.0 | |
| 17 | ||
| 18 | ||
| 19 | def connect(): | |
| 20 | return sqlite3.connect(DB) | |
| 21 | ||
| 22 | ||
| 23 | def load_incidents(conn, agencies=None, start=None, end=None, categories=None): | |
| 24 | where, params = ["lat IS NOT NULL", "lon IS NOT NULL"], [] | |
| 25 | if agencies: | |
| 26 | where.append(f"agency IN ({','.join('?' * len(agencies))})") | |
| 27 | params += list(agencies) | |
| 28 | if categories: | |
| 29 | where.append(f"category IN ({','.join('?' * len(categories))})") | |
| 30 | params += list(categories) | |
| 31 | if start: | |
| 32 | where.append("occurred_at >= ?") | |
| 33 | params.append(f"{start}T00:00:00") | |
| 34 | if end: | |
| 35 | where.append("occurred_at <= ?") | |
| 36 | params.append(f"{end}T23:59:59") | |
| 37 | df = pd.read_sql_query( | |
| 38 | f"""SELECT source, agency, case_id, occurred_at, category, call_type, | |
| 39 | disposition, offense_desc, is_stop, address, lat, lon | |
| 40 | FROM incidents WHERE {' AND '.join(where)}""", | |
| 41 | conn, params=params) | |
| 42 | df["occurred_at"] = pd.to_datetime(df["occurred_at"]) | |
| 43 | return df | |
| 44 | ||
| 45 | ||
| 46 | def load_cameras(conn): | |
| 47 | return pd.read_sql_query( | |
| 48 | "SELECT osm_id, lat, lon, manufacturer, operator, direction, first_seen" | |
| 49 | " FROM alpr_cameras", conn) | |
| 50 | ||
| 51 | ||
| 52 | def nearest_camera_m(df, cameras): | |
| 53 | """Great-circle distance in metres from each incident to the closest camera. | |
| 54 | ||
| 55 | 164 cameras against ~250k incidents, so the full distance matrix is computed | |
| 56 | in chunks rather than all at once.""" | |
| 57 | if df.empty or cameras.empty: | |
| 58 | return np.full(len(df), np.nan) | |
| 59 | ||
| 60 | lat1 = np.radians(df["lat"].to_numpy(dtype=float)) | |
| 61 | lon1 = np.radians(df["lon"].to_numpy(dtype=float)) | |
| 62 | lat2 = np.radians(cameras["lat"].to_numpy(dtype=float)) | |
| 63 | lon2 = np.radians(cameras["lon"].to_numpy(dtype=float)) | |
| 64 | ||
| 65 | out = np.empty(len(df)) | |
| 66 | cos2 = np.cos(lat2) | |
| 67 | for i in range(0, len(df), 20000): | |
| 68 | s = slice(i, i + 20000) | |
| 69 | dlat = lat2[None, :] - lat1[s, None] | |
| 70 | dlon = lon2[None, :] - lon1[s, None] | |
| 71 | a = (np.sin(dlat / 2) ** 2 | |
| 72 | + np.cos(lat1[s, None]) * cos2[None, :] * np.sin(dlon / 2) ** 2) | |
| 73 | out[s] = (2 * EARTH_M * np.arcsin(np.sqrt(a))).min(axis=1) | |
| 74 | return out | |
| 75 | ||
| 76 | ||
| 77 | def daily_counts(df): | |
| 78 | if df.empty: | |
| 79 | return pd.DataFrame(columns=["date", "agency", "incidents"]) | |
| 80 | g = (df.assign(date=df["occurred_at"].dt.floor("D")) | |
| 81 | .groupby(["date", "agency"], as_index=False) | |
| 82 | .size().rename(columns={"size": "incidents"})) | |
| 83 | return g | |
| 84 | ||
| 85 | ||
| 86 | def stop_outcomes(df): | |
| 87 | """Citation and arrest rate on vehicle stops, per agency. | |
| 88 | ||
| 89 | The two CAD systems use different disposition vocabularies -- Sarpy writes | |
| 90 | WRITTEN WARNING / CITATION, Council Bluffs writes "3 - Citation" and folds | |
| 91 | warnings into "7 - Handled by Officer" -- and both allow several outcomes per | |
| 92 | stop. Substring matching on citation and arrest is the only comparison the | |
| 93 | two vocabularies actually support, and a department that records warnings | |
| 94 | less thoroughly will show a higher citation rate for that reason alone.""" | |
| 95 | stops = df[(df["is_stop"] == 1) & df["disposition"].notna()] | |
| 96 | if stops.empty: | |
| 97 | return pd.DataFrame(columns=["agency", "outcome", "rate", "stops"]) | |
| 98 | d = stops["disposition"].str.upper() | |
| 99 | stops = stops.assign(cited=d.str.contains("CITATION"), | |
| 100 | arrested=d.str.contains("ARREST")) | |
| 101 | g = stops.groupby("agency").agg(stops=("cited", "size"), | |
| 102 | Cited=("cited", "mean"), | |
| 103 | Arrested=("arrested", "mean")).reset_index() | |
| 104 | return (g.melt(id_vars=["agency", "stops"], value_vars=["Cited", "Arrested"], | |
| 105 | var_name="outcome", value_name="rate") | |
| 106 | .sort_values("rate", ascending=False)) | |
| 107 | ||
| 108 | ||
| 109 | def camera_proximity(df, cameras, bin_m=200, max_m=2000): | |
| 110 | """Share of stops vs other incidents falling in each distance band. | |
| 111 | ||
| 112 | Both series are normalised, so a gap between them means stops cluster | |
| 113 | differently around cameras than the rest of the call volume does. It is not | |
| 114 | evidence of causation: cameras and stops both concentrate on arterials.""" | |
| 115 | if df.empty or cameras.empty: | |
| 116 | return pd.DataFrame(columns=["distance_m", "kind", "share"]) | |
| 117 | d = df.assign(dist=nearest_camera_m(df, cameras)) | |
| 118 | d = d[d["dist"] <= max_m] | |
| 119 | if d.empty: | |
| 120 | return pd.DataFrame(columns=["distance_m", "kind", "share"]) | |
| 121 | d["kind"] = np.where(d["is_stop"] == 1, "Vehicle stops", "All other incidents") | |
| 122 | d["distance_m"] = (d["dist"] // bin_m * bin_m).astype(int) | |
| 123 | g = (d.groupby(["kind", "distance_m"], as_index=False) | |
| 124 | .size().rename(columns={"size": "n"})) | |
| 125 | g["share"] = g["n"] / g.groupby("kind")["n"].transform("sum") | |
| 126 | return g | |
| 127 | ||
| 128 | ||
| 129 | def agency_options(conn): | |
| 130 | rows = conn.execute( | |
| 131 | "SELECT agency, COUNT(*) FROM incidents GROUP BY agency ORDER BY 2 DESC" | |
| 132 | ).fetchall() | |
| 133 | return [a for a, _ in rows] | |
| 134 | ||
| 135 | ||
| 136 | def category_options(conn): | |
| 137 | rows = conn.execute( | |
| 138 | "SELECT category, COUNT(*) FROM incidents WHERE category IS NOT NULL" | |
| 139 | " GROUP BY category ORDER BY 2 DESC").fetchall() | |
| 140 | return [c for c, _ in rows] | |
| 141 | ||
| 142 | ||
| 143 | def date_bounds(conn): | |
| 144 | lo, hi = conn.execute( | |
| 145 | "SELECT MIN(occurred_at), MAX(occurred_at) FROM incidents").fetchone() | |
| 146 | return lo[:10], hi[:10] | |
app.py +216 −101
| @@ -1,116 +1,231 @@ | ||
| 1 | from dash import Dash, html, dcc, callback, Output, Input, dash_table | |
| 1 | """Omaha metro police activity dashboard. | |
| 2 | ||
| 3 | Covers Omaha PD, Council Bluffs PD and the Sarpy County agencies (Bellevue, | |
| 4 | Papillion, La Vista, Sheriff). Ralston PD publishes no machine-readable feed and | |
| 5 | is absent. OPD publishes NIBRS offence records with no stop or disposition data, | |
| 6 | so it contributes nothing to the enforcement panels. | |
| 7 | """ | |
| 8 | ||
| 9 | from dash import Dash, Input, Output, callback, dash_table, dcc, html | |
| 2 | 10 | import plotly.express as px |
| 3 | import pandas as pd | |
| 4 | import numpy as np | |
| 5 | import sqlite3 | |
| 6 | from datetime import datetime, date | |
| 11 | import plotly.graph_objects as go | |
| 12 | ||
| 13 | import analysis | |
| 14 | ||
| 15 | conn = analysis.connect() | |
| 16 | INCIDENTS = analysis.load_incidents(conn) | |
| 17 | CAMERAS = analysis.load_cameras(conn) | |
| 18 | AGENCIES = analysis.agency_options(conn) | |
| 19 | CATEGORIES = analysis.category_options(conn) | |
| 20 | DATE_LO, DATE_HI = analysis.date_bounds(conn) | |
| 21 | conn.close() | |
| 7 | 22 | |
| 8 | # Connect to database and query all incidents | |
| 9 | connection = sqlite3.connect("./raw_data/ingress.db") | |
| 10 | cursor = connection.cursor() | |
| 11 | query = "SELECT * FROM incidents;" | |
| 12 | df = pd.read_sql_query(query, connection).sort_values(by="description") | |
| 23 | DEFAULT_AGENCIES = [a for a in analysis.POLICE_AGENCIES if a in AGENCIES] | |
| 24 | CENTER = {"lat": 41.21, "lon": -95.97} | |
| 25 | MAP_SAMPLE = 15000 | |
| 26 | # Plotly writes colours into the figure, so the scheme has to be known before a | |
| 27 | # figure is built. assets/theme.js reports the media query into the theme store | |
| 28 | # and every figure callback reads it; nothing is restyled after the fact. | |
| 29 | PALETTES = { | |
| 30 | "light": {"fg": "#111", "grid": "#e6e6e6", "legend": "rgba(255,255,255,.85)", | |
| 31 | "basemap": "open-street-map"}, | |
| 32 | "dark": {"fg": "#e8e8ea", "grid": "#333840", "legend": "rgba(27,30,36,.85)", | |
| 33 | "basemap": "carto-darkmatter"}, | |
| 34 | } | |
| 35 | ||
| 36 | ||
| 37 | def themed(fig, theme): | |
| 38 | p = PALETTES.get(theme, PALETTES["light"]) | |
| 39 | fig.update_layout(paper_bgcolor="rgba(0,0,0,0)", plot_bgcolor="rgba(0,0,0,0)", | |
| 40 | font_color=p["fg"], legend_bgcolor=p["legend"], | |
| 41 | legend_font_color=p["fg"]) | |
| 42 | # update_xaxes would bolt empty axis objects onto the map figure, which has | |
| 43 | # no cartesian axes at all. | |
| 44 | if not any(tr.type == "scattermap" for tr in fig.data): | |
| 45 | fig.update_xaxes(gridcolor=p["grid"], zerolinecolor=p["grid"]) | |
| 46 | fig.update_yaxes(gridcolor=p["grid"], zerolinecolor=p["grid"]) | |
| 47 | return fig | |
| 13 | 48 | |
| 14 | # Replace empty cells with NaN | |
| 15 | df = df.replace(r'^\s*$', np.nan, regex=True) | |
| 16 | 49 | |
| 17 | # Convert date column to datetime | |
| 18 | # TODO: Create a combined datetime column in the db to load/convert here | |
| 19 | df["date"] = pd.to_datetime(df["date"]) | |
| 50 | def empty(theme): | |
| 51 | return themed(go.Figure().update_layout( | |
| 52 | annotations=[{"text": "No incidents match these filters", | |
| 53 | "showarrow": False, "font": {"size": 14}}], | |
| 54 | xaxis={"visible": False}, yaxis={"visible": False}), theme) | |
| 20 | 55 | |
| 21 | # Configure HTML layout | |
| 22 | 56 | app = Dash(__name__) |
| 23 | app.layout = html.Div(children = [ | |
| 24 | html.H1(children="Omaha Crime Mapping & Analysis", style={"textAlign":"center"}), | |
| 25 | html.Div([ | |
| 26 | html.Div(className="flex", | |
| 27 | children = [ | |
| 28 | html.P("Crime:"), | |
| 29 | dcc.Dropdown(df.sort_values("description").description.unique(), "INJURY", id="dropdown"), | |
| 30 | ] | |
| 31 | ), | |
| 32 | html.Div(className="flex", | |
| 33 | children = [ | |
| 34 | html.P("Date Range:"), | |
| 35 | dcc.DatePickerRange( | |
| 36 | id='date-picker-range', | |
| 37 | min_date_allowed=date(2015, 1, 1), | |
| 38 | max_date_allowed=date(2023, 12, 31), | |
| 39 | start_date=date(2015, 1, 1), | |
| 40 | end_date=date(2023,12,31) | |
| 41 | ), | |
| 42 | ] | |
| 43 | ), | |
| 44 | dcc.Graph(id="map-graph"), | |
| 45 | dcc.Graph(id="line-graph"), | |
| 46 | dash_table.DataTable(data=df.to_dict("records"), page_size=5, id="table") | |
| 47 | ]) | |
| 57 | app.title = "Omaha Metro Police Activity" | |
| 58 | ||
| 59 | app.layout = html.Div([ | |
| 60 | dcc.Store(id="theme", data="light"), | |
| 61 | html.H1("Omaha Metro Police Activity"), | |
| 62 | html.P([ | |
| 63 | f"{len(INCIDENTS):,} geolocated incidents, {DATE_LO} to {DATE_HI}. ", | |
| 64 | f"{len(CAMERAS)} ALPR cameras from OpenStreetMap. ", | |
| 65 | "Ralston PD publishes no feed and is not represented. ", | |
| 66 | "Omaha PD reports no stops or dispositions, so it is absent from the ", | |
| 67 | "enforcement panels below.", | |
| 68 | ], className="subtitle"), | |
| 69 | ||
| 70 | html.Div(className="controls", children=[ | |
| 71 | html.Div([html.Label("Agency"), | |
| 72 | dcc.Dropdown(AGENCIES, DEFAULT_AGENCIES, id="agencies", | |
| 73 | multi=True)]), | |
| 74 | html.Div([html.Label("Category"), | |
| 75 | dcc.Dropdown(CATEGORIES, [], id="categories", multi=True, | |
| 76 | placeholder="all categories")]), | |
| 77 | html.Div([html.Label("Dates"), | |
| 78 | dcc.DatePickerRange(id="dates", min_date_allowed=DATE_LO, | |
| 79 | max_date_allowed=DATE_HI, | |
| 80 | start_date=DATE_LO, end_date=DATE_HI)]), | |
| 81 | html.Div([html.Label("Show cameras"), | |
| 82 | dcc.Checklist([{"label": " ALPR layer", "value": "on"}], | |
| 83 | ["on"], id="show-cameras")]), | |
| 84 | ]), | |
| 85 | ||
| 86 | html.Div(id="kpis", className="kpis"), | |
| 87 | ||
| 88 | dcc.Graph(id="map"), | |
| 89 | dcc.Graph(id="timeline"), | |
| 90 | html.Div(className="row", children=[ | |
| 91 | dcc.Graph(id="dispositions", className="half"), | |
| 92 | dcc.Graph(id="proximity", className="half"), | |
| 93 | ]), | |
| 94 | html.H2("Incidents"), | |
| 95 | dash_table.DataTable(id="table", page_size=10, sort_action="native", | |
| 96 | style_table={"overflowX": "auto"}), | |
| 48 | 97 | ]) |
| 49 | 98 | |
| 50 | # Create map | |
| 51 | @callback( | |
| 52 | Output("map-graph", "figure"), | |
| 53 | Input("dropdown", "value"), | |
| 54 | Input('date-picker-range', 'start_date'), | |
| 55 | Input('date-picker-range', 'end_date') | |
| 56 | ) | |
| 57 | def update_map(description, start_date, end_date): | |
| 58 | dff = df[(df['date'] > start_date) & (df['date'] < end_date)] | |
| 59 | dff = dff.reset_index() | |
| 60 | dff = dff[dff.description == description] | |
| 61 | ||
| 62 | fig = px.scatter_mapbox( | |
| 63 | dff, | |
| 64 | lat="lat", | |
| 65 | lon="lon", | |
| 66 | hover_name="description", | |
| 67 | hover_data=["date", "time"], | |
| 68 | title="Incident Count by Coordinates", | |
| 69 | center={"lat": 41.257160, "lon": -95.995102}, | |
| 70 | zoom=10 | |
| 71 | ) | |
| 72 | ||
| 73 | fig.update_layout(showlegend=False) | |
| 74 | fig.update_layout(mapbox_style="open-street-map") | |
| 75 | fig.update_layout(margin={"r": 0, "t": 0, "l": 0, "b": 0}) | |
| 76 | fig.update_layout(mapbox_bounds={"west": -180, "east": -50, "south": 20, "north": 90}) | |
| 77 | 99 | |
| 78 | return fig | |
| 100 | def filtered(agencies, categories, start, end): | |
| 101 | df = INCIDENTS | |
| 102 | if agencies: | |
| 103 | df = df[df["agency"].isin(agencies)] | |
| 104 | if categories: | |
| 105 | df = df[df["category"].isin(categories)] | |
| 106 | if start: | |
| 107 | df = df[df["occurred_at"] >= start] | |
| 108 | if end: | |
| 109 | df = df[df["occurred_at"] <= f"{end[:10]} 23:59:59"] | |
| 110 | return df | |
| 79 | 111 | |
| 80 | # Create line graph | |
| 81 | @callback( | |
| 82 | Output("line-graph", "figure"), | |
| 83 | Input("dropdown", "value"), | |
| 84 | Input('date-picker-range', 'start_date'), | |
| 85 | Input('date-picker-range', 'end_date') | |
| 86 | ) | |
| 87 | def update_line(description, start_date, end_date): | |
| 88 | dff = df[(df['date'] > start_date) & (df['date'] < end_date)] | |
| 89 | dff = dff.reset_index() | |
| 90 | dff = dff[dff.description == description] | |
| 91 | dff = dff.groupby(by="date").count() | |
| 92 | dff = dff.reset_index() | |
| 93 | ||
| 94 | fig = px.line( | |
| 95 | dff, | |
| 96 | x="date", | |
| 97 | y="description", | |
| 98 | ) | |
| 99 | 112 | |
| 100 | return fig | |
| 113 | INPUTS = [Input("agencies", "value"), Input("categories", "value"), | |
| 114 | Input("dates", "start_date"), Input("dates", "end_date"), | |
| 115 | Input("theme", "data")] | |
| 116 | ||
| 117 | ||
| 118 | @callback(Output("kpis", "children"), *INPUTS) | |
| 119 | def update_kpis(agencies, categories, start, end, _theme): | |
| 120 | df = filtered(agencies, categories, start, end) | |
| 121 | stops = df[df["is_stop"] == 1] | |
| 122 | # Rate over the span the stops actually cover, not the full incident range: | |
| 123 | # OPD reaches back to 2022 but reports no stops at all. | |
| 124 | if len(stops): | |
| 125 | span = (stops["occurred_at"].max() - stops["occurred_at"].min()).days | |
| 126 | rate = f"{len(stops) / max(span, 1):.1f}" | |
| 127 | else: | |
| 128 | rate = "0" | |
| 129 | cards = [ | |
| 130 | ("Incidents", f"{len(df):,}"), | |
| 131 | ("Vehicle stops", f"{len(stops):,}"), | |
| 132 | ("Stops per day", rate), | |
| 133 | ("Agencies", f"{df['agency'].nunique()}"), | |
| 134 | ] | |
| 135 | return [html.Div(className="kpi", children=[html.Span(v, className="kpi-value"), | |
| 136 | html.Span(k, className="kpi-label")]) | |
| 137 | for k, v in cards] | |
| 138 | ||
| 139 | ||
| 140 | @callback(Output("map", "figure"), *INPUTS, Input("show-cameras", "value")) | |
| 141 | def update_map(agencies, categories, start, end, theme, show_cameras): | |
| 142 | df = filtered(agencies, categories, start, end) | |
| 143 | # A density layer over a quarter-million points saturates at any radius and | |
| 144 | # hides which agency is where, so plot a sample of the points instead. | |
| 145 | sampled = len(df) > MAP_SAMPLE | |
| 146 | if sampled: | |
| 147 | df = df.sample(MAP_SAMPLE, random_state=0) | |
| 148 | ||
| 149 | fig = go.Figure() | |
| 150 | for agency, g in df.groupby("agency", sort=False): | |
| 151 | fig.add_trace(go.Scattermap( | |
| 152 | lat=g["lat"], lon=g["lon"], mode="markers", name=agency, | |
| 153 | marker={"size": 4, "opacity": 0.45}, | |
| 154 | text=g["call_type"].fillna(g["category"]).fillna(g["offense_desc"]), | |
| 155 | hovertemplate="%{text}<extra>" + agency + "</extra>")) | |
| 156 | if show_cameras and len(CAMERAS): | |
| 157 | fig.add_trace(go.Scattermap( | |
| 158 | lat=CAMERAS["lat"], lon=CAMERAS["lon"], mode="markers", | |
| 159 | # Amber reads against both the light and the dark basemap. | |
| 160 | marker={"size": 7, "color": "#ffb300"}, | |
| 161 | name="ALPR camera", | |
| 162 | text=[f"{m or 'unknown make'} — {o or 'operator not tagged'}" | |
| 163 | for m, o in zip(CAMERAS["manufacturer"], CAMERAS["operator"])], | |
| 164 | hovertemplate="%{text}<extra>ALPR</extra>")) | |
| 165 | title = f"{MAP_SAMPLE:,}-incident sample" if sampled else f"{len(df):,} incidents" | |
| 166 | fig.update_layout(map={"style": PALETTES[theme]["basemap"], "center": CENTER, | |
| 167 | "zoom": 9.6}, | |
| 168 | height=560, margin={"r": 0, "t": 30, "l": 0, "b": 0}, | |
| 169 | title=title, uirevision="map", | |
| 170 | legend={"x": 0.01, "y": 0.99}) | |
| 171 | return themed(fig, theme) | |
| 172 | ||
| 173 | ||
| 174 | @callback(Output("timeline", "figure"), *INPUTS) | |
| 175 | def update_timeline(agencies, categories, start, end, theme): | |
| 176 | df = filtered(agencies, categories, start, end) | |
| 177 | counts = analysis.daily_counts(df) | |
| 178 | if counts.empty: | |
| 179 | return empty(theme) | |
| 180 | fig = px.line(counts, x="date", y="incidents", color="agency", | |
| 181 | title="Incidents per day by agency") | |
| 182 | fig.update_layout(margin={"t": 40}, hovermode="x unified") | |
| 183 | return themed(fig, theme) | |
| 184 | ||
| 185 | ||
| 186 | @callback(Output("dispositions", "figure"), *INPUTS) | |
| 187 | def update_dispositions(agencies, categories, start, end, theme): | |
| 188 | # Category filter is ignored: stops are identified by is_stop, not category. | |
| 189 | out = analysis.stop_outcomes(filtered(agencies, None, start, end)) | |
| 190 | if out.empty: | |
| 191 | return empty(theme) | |
| 192 | fig = px.bar(out, x="rate", y="agency", color="outcome", orientation="h", | |
| 193 | barmode="group", custom_data=["stops"], | |
| 194 | title="Vehicle stop outcomes", | |
| 195 | labels={"rate": "share of that agency's stops"}) | |
| 196 | fig.update_traces(hovertemplate="%{x:.1%} of %{customdata[0]:,} stops" | |
| 197 | "<extra>%{fullData.name}</extra>") | |
| 198 | fig.update_layout(margin={"t": 40}, xaxis_tickformat=".0%", | |
| 199 | yaxis_title=None, legend_title_text=None) | |
| 200 | return themed(fig, theme) | |
| 201 | ||
| 202 | ||
| 203 | @callback(Output("proximity", "figure"), *INPUTS) | |
| 204 | def update_proximity(agencies, categories, start, end, theme): | |
| 205 | # Category filter is deliberately ignored: the comparison needs both stops | |
| 206 | # and the non-stop baseline in the same window. | |
| 207 | df = filtered(agencies, None, start, end) | |
| 208 | prox = analysis.camera_proximity(df, CAMERAS) | |
| 209 | if prox.empty: | |
| 210 | return empty(theme) | |
| 211 | fig = px.line(prox, x="distance_m", y="share", color="kind", markers=True, | |
| 212 | title="Distance to nearest ALPR camera", | |
| 213 | labels={"distance_m": "metres to nearest camera", | |
| 214 | "share": "share of incidents"}) | |
| 215 | fig.update_layout(margin={"t": 40}, yaxis_tickformat=".1%") | |
| 216 | return themed(fig, theme) | |
| 217 | ||
| 218 | ||
| 219 | @callback(Output("table", "data"), Output("table", "columns"), *INPUTS) | |
| 220 | def update_table(agencies, categories, start, end, _theme): | |
| 221 | cols = ["occurred_at", "agency", "category", "call_type", "disposition", | |
| 222 | "address"] | |
| 223 | df = (filtered(agencies, categories, start, end) | |
| 224 | .sort_values("occurred_at", ascending=False) | |
| 225 | .head(500)[cols]) | |
| 226 | df = df.assign(occurred_at=df["occurred_at"].dt.strftime("%Y-%m-%d %H:%M")) | |
| 227 | return df.to_dict("records"), [{"name": c, "id": c} for c in cols] | |
| 101 | 228 | |
| 102 | # Create table | |
| 103 | @callback( | |
| 104 | Output("table", "data"), | |
| 105 | Input("dropdown", "value"), | |
| 106 | Input('date-picker-range', 'start_date'), | |
| 107 | Input('date-picker-range', 'end_date') | |
| 108 | ) | |
| 109 | def update_table(description, start_date, end_date): | |
| 110 | dff = df[(df['date'] > start_date) & (df['date'] < end_date)] | |
| 111 | dff = dff.reset_index() | |
| 112 | dff = dff[dff.description == description] | |
| 113 | return dff.to_dict("records") | |
| 114 | 229 | |
| 115 | 230 | if __name__ == "__main__": |
| 116 | 231 | app.run(debug=True) |
assets/styles.css +96 −11
| @@ -1,4 +1,37 @@ | ||
| 1 | :root { | |
| 2 | --bg: #fff; | |
| 3 | --fg: #111; | |
| 4 | --muted: #555; | |
| 5 | --border: #ddd; | |
| 6 | --panel: #fafafa; | |
| 7 | } | |
| 8 | ||
| 9 | @media (prefers-color-scheme: dark) { | |
| 10 | :root { | |
| 11 | --bg: #14161a; | |
| 12 | --fg: #e8e8ea; | |
| 13 | --muted: #9aa0a6; | |
| 14 | --border: #333840; | |
| 15 | --panel: #1b1e24; | |
| 16 | ||
| 17 | /* Dash 3 ships light-only values for its component design tokens. */ | |
| 18 | --Dash-Text-Primary: #e8e8ea; | |
| 19 | --Dash-Text-Strong: #f2f2f4; | |
| 20 | --Dash-Text-Weak: #9aa0a6; | |
| 21 | --Dash-Text-Disabled: #6b7178; | |
| 22 | --Dash-Stroke-Strong: rgba(232, 232, 234, 0.45); | |
| 23 | --Dash-Stroke-Weak: rgba(232, 232, 234, 0.15); | |
| 24 | --Dash-Fill-Inverse-Strong: #1b1e24; | |
| 25 | --Dash-Fill-Interactive-Weak: rgba(255, 255, 255, 0.06); | |
| 26 | --Dash-Fill-Primary-Hover: rgba(255, 255, 255, 0.06); | |
| 27 | --Dash-Fill-Primary-Active: rgba(255, 255, 255, 0.1); | |
| 28 | --Dash-Fill-Disabled: rgba(255, 255, 255, 0.1); | |
| 29 | } | |
| 30 | } | |
| 31 | ||
| 1 | 32 | body { |
| 33 | background: var(--bg); | |
| 34 | color: var(--fg); | |
| 2 | 35 | font-family: -apple-system, BlinkMacSystemFont, avenir next, avenir, segoe ui, helvetica neue, helvetica, Cantarell, Ubuntu, roboto, noto, arial, sans-serif; |
| 3 | 36 | max-width: 80vw; |
| 4 | 37 | margin: 2rem auto; |
| @@ -14,21 +47,73 @@ body { | ||
| 14 | 47 | max-width: 100%; |
| 15 | 48 | } |
| 16 | 49 | |
| 17 | .flex { | |
| 18 | display: flex; | |
| 19 | flex-direction: row; | |
| 20 | align-items: center; | |
| 50 | .subtitle { | |
| 51 | color: var(--muted); | |
| 52 | margin-top: -0.5rem; | |
| 53 | } | |
| 54 | ||
| 55 | .controls { | |
| 56 | display: grid; | |
| 57 | grid-template-columns: repeat(auto-fit, minmax(220px, 1fr)); | |
| 58 | gap: 1rem; | |
| 59 | margin: 1.5rem 0; | |
| 60 | } | |
| 61 | ||
| 62 | .controls label { | |
| 63 | display: block; | |
| 64 | font-size: 0.85rem; | |
| 65 | font-weight: 600; | |
| 66 | margin-bottom: 0.25rem; | |
| 21 | 67 | } |
| 22 | 68 | |
| 23 | .flex p { | |
| 24 | margin-right: 1rem; | |
| 25 | flex-shrink: 0; | |
| 69 | .kpis { | |
| 70 | display: grid; | |
| 71 | grid-template-columns: repeat(auto-fit, minmax(150px, 1fr)); | |
| 72 | gap: 1rem; | |
| 26 | 73 | } |
| 27 | 74 | |
| 28 | .flex .dash-dropdown { | |
| 29 | width: 100%; | |
| 75 | .kpi { | |
| 76 | background: var(--panel); | |
| 77 | border: 1px solid var(--border); | |
| 78 | border-radius: 6px; | |
| 79 | padding: 0.75rem 1rem; | |
| 30 | 80 | } |
| 31 | 81 | |
| 32 | .flex #date-picker-range { | |
| 33 | width: 100%; | |
| 82 | .kpi-value { | |
| 83 | display: block; | |
| 84 | font-size: 1.6rem; | |
| 85 | font-weight: 600; | |
| 86 | } | |
| 87 | ||
| 88 | .kpi-label { | |
| 89 | display: block; | |
| 90 | font-size: 0.8rem; | |
| 91 | color: var(--muted); | |
| 92 | } | |
| 93 | ||
| 94 | .row { | |
| 95 | display: grid; | |
| 96 | grid-template-columns: repeat(auto-fit, minmax(420px, 1fr)); | |
| 97 | gap: 1rem; | |
| 98 | } | |
| 99 | ||
| 100 | .row .half { | |
| 101 | min-width: 0; | |
| 102 | } | |
| 103 | ||
| 104 | /* DataTable paints its own colours instead of reading the tokens above. */ | |
| 105 | @media (prefers-color-scheme: dark) { | |
| 106 | .dash-spreadsheet-inner td, | |
| 107 | .dash-spreadsheet-inner th, | |
| 108 | .dash-spreadsheet-inner td.dash-cell, | |
| 109 | .dash-spreadsheet-inner th.dash-header { | |
| 110 | background-color: var(--panel) !important; | |
| 111 | border-color: var(--border) !important; | |
| 112 | color: var(--fg) !important; | |
| 113 | } | |
| 114 | ||
| 115 | .dash-spreadsheet-menu, | |
| 116 | .dash-spreadsheet-container { | |
| 117 | color: var(--fg); | |
| 118 | } | |
| 34 | 119 | } |
assets/theme.js added +19
| @@ -0,0 +1,19 @@ | ||
| 1 | // Plotly bakes colours into the rendered figure, so the server needs to know the | |
| 2 | // colour scheme before it builds one. Push the media query result into the | |
| 3 | // theme store on load and whenever it changes; app.py styles from there. | |
| 4 | (function () { | |
| 5 | var query = window.matchMedia('(prefers-color-scheme: dark)'); | |
| 6 | ||
| 7 | function publish(tries) { | |
| 8 | var dc = window.dash_clientside; | |
| 9 | try { | |
| 10 | // set_props needs the layout rendered, not just the bundle loaded. | |
| 11 | dc.set_props('theme', {data: query.matches ? 'dark' : 'light'}); | |
| 12 | } catch (e) { | |
| 13 | if ((tries || 0) < 100) setTimeout(function () { publish((tries || 0) + 1); }, 50); | |
| 14 | } | |
| 15 | } | |
| 16 | ||
| 17 | query.addEventListener('change', function () { publish(0); }); | |
| 18 | publish(0); | |
| 19 | })(); | |
ingest.py +332 −75
| @@ -1,77 +1,334 @@ | ||
| 1 | # Import required modules | |
| 1 | """Pull incident feeds and ALPR camera locations into raw_data/metro.db. | |
| 2 | ||
| 3 | Sources: | |
| 4 | opd Omaha Police incident data (DCGIS ArcGIS view), 2022-01-01 onward, daily. | |
| 5 | sarpy Sarpy County PublicCrimeMap CAD calls, rolling 12-month window. Covers | |
| 6 | Bellevue, Papillion, La Vista, Gretna, Springfield and the Sheriff. | |
| 7 | Records age out of the feed, so run this often enough to keep the archive. | |
| 8 | cbpd Council Bluffs PD public CFS feed, rolling 12-month window, refreshed | |
| 9 | every 10 minutes. Same ageing-out caveat as sarpy. | |
| 10 | alpr ALPR cameras from OpenStreetMap via Overpass (the data behind DeFlock). | |
| 11 | opd_csv One-time backfill of raw_data/Incidents_*.csv (2015-2023). Statute text | |
| 12 | only, no NIBRS category. | |
| 13 | ||
| 14 | All three ArcGIS services return UTC epochs. Their WHERE literals do not agree: | |
| 15 | OPD and Council Bluffs use UTC, Sarpy uses America/Chicago, so each source | |
| 16 | carries its own literal timezone and page size. | |
| 17 | """ | |
| 18 | ||
| 19 | import argparse | |
| 2 | 20 | import csv |
| 21 | import json | |
| 3 | 22 | import sqlite3 |
| 4 | import os | |
| 5 | ||
| 6 | # Create the database file | |
| 7 | connection = sqlite3.connect('./raw_data/test.db') | |
| 8 | ||
| 9 | # Creating a cursor object to execute SQL queries | |
| 10 | cursor = connection.cursor() | |
| 11 | ||
| 12 | # Table Definition | |
| 13 | # rb = RB Number | |
| 14 | # date = Reported Date | |
| 15 | # time = Reported Time | |
| 16 | # description = Statute/Ordinance Description | |
| 17 | # location = Occurred Location | |
| 18 | # district = Occurred District | |
| 19 | # lat = Occurred Block LAT | |
| 20 | # lon = Occurred Block LON | |
| 21 | create_table = '''CREATE TABLE incidents( | |
| 22 | id INTEGER PRIMARY KEY AUTOINCREMENT, | |
| 23 | rb TEXT NOT NULL, | |
| 24 | date TEXT NOT NULL, | |
| 25 | time TEXT NOT NULL, | |
| 26 | description TEXT NOT NULL, | |
| 27 | location TEXT NOT NULL, | |
| 28 | district TEXT NOT NULL, | |
| 29 | lat REAL NOT NULL, | |
| 30 | lon REAL NOT NULL); | |
| 31 | ''' | |
| 32 | ||
| 33 | # Create the table | |
| 34 | cursor.execute(create_table) | |
| 35 | ||
| 36 | # Point to the data directory | |
| 37 | directory = os.fsencode("./raw_data/") | |
| 38 | ||
| 39 | # Loop through all raw data files | |
| 40 | for file in os.listdir(directory): | |
| 41 | filename = os.fsdecode(file) | |
| 42 | if filename.endswith(".csv"): | |
| 43 | # Opening the file | |
| 44 | file = open("./raw_data/" + filename) | |
| 45 | ||
| 46 | # Reading the contents of the file | |
| 47 | contents = csv.reader(file) | |
| 48 | ||
| 49 | # SQL query to insert data into the | |
| 50 | # table | |
| 51 | insert_records = "INSERT INTO incidents (rb, date, time, description, location, district, lat, lon) VALUES(?, ?, ?, ?, ?, ?, ?, ?)" | |
| 52 | ||
| 53 | # Importing the contents of the file | |
| 54 | # into our table | |
| 55 | cursor.executemany(insert_records, contents) | |
| 56 | print("Inserted data from: ", filename) | |
| 57 | continue | |
| 58 | else: | |
| 59 | continue | |
| 60 | ||
| 61 | # Delete extra copies of the header row that were inserted | |
| 62 | delete_headers = "DELETE FROM incidents WHERE rb = 'RB Number'" | |
| 63 | cursor.execute(delete_headers) | |
| 64 | ||
| 65 | # Test query to see if the data loaded | |
| 66 | select_all = "SELECT * FROM incidents" | |
| 67 | rows = cursor.execute(select_all).fetchall() | |
| 68 | ||
| 69 | # Output to the console screen | |
| 70 | for r in rows: | |
| 71 | print(r) | |
| 72 | ||
| 73 | # Commit the changes | |
| 74 | connection.commit() | |
| 75 | ||
| 76 | # Close the database connection | |
| 77 | connection.close() | |
| 23 | import ssl | |
| 24 | import sys | |
| 25 | import time | |
| 26 | import urllib.error | |
| 27 | import urllib.parse | |
| 28 | import urllib.request | |
| 29 | ||
| 30 | import certifi | |
| 31 | from datetime import datetime, timedelta, timezone | |
| 32 | from pathlib import Path | |
| 33 | from zoneinfo import ZoneInfo | |
| 34 | ||
| 35 | # Python does not use the macOS keychain, and Council Bluffs' host is not in the | |
| 36 | # bundled trust store on every platform. | |
| 37 | SSL_CTX = ssl.create_default_context(cafile=certifi.where()) | |
| 38 | ||
| 39 | ROOT = Path(__file__).parent | |
| 40 | DB = ROOT / "raw_data" / "metro.db" | |
| 41 | LOCAL = ZoneInfo("America/Chicago") | |
| 42 | ||
| 43 | OPD = { | |
| 44 | "url": "https://services1.arcgis.com/tIBLyYZX96jUntYm/arcgis/rest/services" | |
| 45 | "/Omaha_Police_Incident_Data_(View)/FeatureServer/0", | |
| 46 | "date_field": "dteMidpoint", | |
| 47 | "literal_tz": timezone.utc, | |
| 48 | "oid": "OBJECTID", | |
| 49 | "page": 2000, | |
| 50 | } | |
| 51 | ||
| 52 | SARPY = { | |
| 53 | "url": "https://geodata.sarpy.gov/arcgis/rest/services/PublicSafety" | |
| 54 | "/PublicCrimeMap/FeatureServer/1", | |
| 55 | "date_field": "IncidentDate", | |
| 56 | "literal_tz": LOCAL, | |
| 57 | "oid": "ObjectID", | |
| 58 | "page": 2000, | |
| 59 | } | |
| 60 | ||
| 61 | CBPD = { | |
| 62 | "url": "https://gispublic.councilbluffs-ia.gov/publicserver/rest/services" | |
| 63 | "/Hosted/Public_Facing_CFS_view/FeatureServer/0", | |
| 64 | "date_field": "cfs_datetime", | |
| 65 | "literal_tz": timezone.utc, | |
| 66 | "oid": "objectid", | |
| 67 | "page": 1000, | |
| 68 | } | |
| 69 | ||
| 70 | # Council Bluffs files officer-initiated stops under one incident_code. | |
| 71 | CBPD_STOP_CODE = "TRAFFIC : TRAFFIC STOP" | |
| 72 | ||
| 73 | # Sarpy IncidentId prefixes. Fire/EMS agencies share the feed with the police | |
| 74 | # agencies; they are kept so the archive stays complete and filtered in the app. | |
| 75 | SARPY_AGENCIES = { | |
| 76 | "LBP": "Bellevue PD", | |
| 77 | "LPP": "Papillion PD", | |
| 78 | "LLP": "La Vista PD", | |
| 79 | "LSO": "Sarpy County SO", | |
| 80 | "LGP": "Gretna PD", | |
| 81 | "LSP": "Springfield PD", | |
| 82 | "BVF": "Bellevue Fire", | |
| 83 | "PAF": "Papillion Fire", | |
| 84 | "GRF": "Gretna Fire", | |
| 85 | "SPF": "Springfield Fire", | |
| 86 | "LVF": "La Vista Fire", | |
| 87 | } | |
| 88 | ||
| 89 | # The feed's other vehicle category, "Traffic", is crashes, parking and DUI | |
| 90 | # calls -- reactive, not officer-initiated. | |
| 91 | SARPY_STOP_CATEGORY = "Proactive Policing - Vehicle Stop" | |
| 92 | ||
| 93 | OVERPASS = "https://overpass-api.de/api/interpreter" | |
| 94 | # Douglas and Sarpy counties in Nebraska plus Council Bluffs across the river. | |
| 95 | BBOX = (40.95, -96.35, 41.45, -95.65) | |
| 96 | OVERPASS_QUERY = f""" | |
| 97 | [out:json][timeout:120]; | |
| 98 | ( | |
| 99 | node["man_made"="surveillance"]["surveillance:type"="ALPR"]{BBOX}; | |
| 100 | node["man_made"="surveillance"]["surveillance:zone"="traffic"]["brand"~"Flock",i]{BBOX}; | |
| 101 | ); | |
| 102 | out body; | |
| 103 | """ | |
| 104 | ||
| 105 | ||
| 106 | def get(url, params, retries=4): | |
| 107 | body = urllib.parse.urlencode(params).encode() | |
| 108 | for attempt in range(retries): | |
| 109 | try: | |
| 110 | req = urllib.request.Request(url, data=body, | |
| 111 | headers={"User-Agent": "omaha-incidents/1.0"}) | |
| 112 | with urllib.request.urlopen(req, timeout=120, context=SSL_CTX) as r: | |
| 113 | payload = json.load(r) | |
| 114 | except (urllib.error.URLError, TimeoutError, json.JSONDecodeError) as e: | |
| 115 | if attempt == retries - 1: | |
| 116 | raise | |
| 117 | time.sleep(2 ** attempt) | |
| 118 | continue | |
| 119 | if "error" in payload: | |
| 120 | raise RuntimeError(f"{url}: {payload['error']}") | |
| 121 | return payload | |
| 122 | raise AssertionError("unreachable") | |
| 123 | ||
| 124 | ||
| 125 | def query_all(service, since): | |
| 126 | """Page through an ArcGIS layer, yielding feature attribute dicts.""" | |
| 127 | where = "1=1" | |
| 128 | if since is not None: | |
| 129 | stamp = since.astimezone(service["literal_tz"]).strftime("%Y-%m-%d %H:%M:%S") | |
| 130 | where = f"{service['date_field']} >= TIMESTAMP '{stamp}'" | |
| 131 | offset = 0 | |
| 132 | while True: | |
| 133 | page = get(service["url"] + "/query", { | |
| 134 | "where": where, | |
| 135 | "outFields": "*", | |
| 136 | "returnGeometry": "true", | |
| 137 | "outSR": "4326", | |
| 138 | "orderByFields": f"{service['oid']} ASC", | |
| 139 | "resultOffset": offset, | |
| 140 | "resultRecordCount": service["page"], | |
| 141 | "f": "json", | |
| 142 | }) | |
| 143 | feats = page.get("features", []) | |
| 144 | if not feats: | |
| 145 | return | |
| 146 | for f in feats: | |
| 147 | yield f | |
| 148 | offset += len(feats) | |
| 149 | print(f" {offset} rows", end="\r", file=sys.stderr, flush=True) | |
| 150 | if not page.get("exceededTransferLimit") and len(feats) < service["page"]: | |
| 151 | return | |
| 152 | ||
| 153 | ||
| 154 | def local_iso(epoch_ms): | |
| 155 | if epoch_ms is None: | |
| 156 | return None | |
| 157 | dt = datetime.fromtimestamp(epoch_ms / 1000, timezone.utc) | |
| 158 | return dt.astimezone(LOCAL).strftime("%Y-%m-%dT%H:%M:%S") | |
| 159 | ||
| 160 | ||
| 161 | def upsert(conn, rows): | |
| 162 | conn.executemany( | |
| 163 | """INSERT INTO incidents | |
| 164 | (source, source_key, agency, case_id, occurred_at, category, | |
| 165 | call_type, disposition, offense_desc, is_stop, address, lat, lon) | |
| 166 | VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?) | |
| 167 | ON CONFLICT(source, source_key) DO UPDATE SET | |
| 168 | agency=excluded.agency, case_id=excluded.case_id, | |
| 169 | occurred_at=excluded.occurred_at, category=excluded.category, | |
| 170 | call_type=excluded.call_type, disposition=excluded.disposition, | |
| 171 | offense_desc=excluded.offense_desc, is_stop=excluded.is_stop, | |
| 172 | address=excluded.address, lat=excluded.lat, lon=excluded.lon""", | |
| 173 | rows) | |
| 174 | return len(rows) | |
| 175 | ||
| 176 | ||
| 177 | def ingest_opd(conn, since): | |
| 178 | rows = [] | |
| 179 | for f in query_all(OPD, since): | |
| 180 | a = f["attributes"] | |
| 181 | occurred = local_iso(a["dteMidpoint"]) | |
| 182 | if occurred is None: | |
| 183 | continue | |
| 184 | rows.append(("opd", str(a["PK"]), "Omaha PD", a.get("RB"), occurred, | |
| 185 | a.get("NIBRSCategory"), None, None, None, 0, | |
| 186 | a.get("AddressBlock"), a.get("LatBlock"), a.get("LonBlock"))) | |
| 187 | return upsert(conn, rows) | |
| 188 | ||
| 189 | ||
| 190 | def ingest_sarpy(conn, since): | |
| 191 | rows, unmapped = [], set() | |
| 192 | for f in query_all(SARPY, since): | |
| 193 | a = f["attributes"] | |
| 194 | occurred = local_iso(a["IncidentDate"]) | |
| 195 | if occurred is None: | |
| 196 | continue | |
| 197 | iid = a["IncidentId"] | |
| 198 | prefix = iid[:3] | |
| 199 | agency = SARPY_AGENCIES.get(prefix) | |
| 200 | if agency is None: | |
| 201 | unmapped.add(prefix) | |
| 202 | agency = f"Unmapped {prefix}" | |
| 203 | g = f.get("geometry") or {} | |
| 204 | rows.append(("sarpy", iid, agency, iid, occurred, a.get("Category"), | |
| 205 | a.get("CadTypeDesc"), a.get("CadDisposition"), | |
| 206 | a.get("StatuteDesc"), | |
| 207 | int(a.get("Category") == SARPY_STOP_CATEGORY), | |
| 208 | a.get("BlkAddress"), g.get("y"), g.get("x"))) | |
| 209 | n = upsert(conn, rows) | |
| 210 | if unmapped: | |
| 211 | print(f" unmapped IncidentId prefixes: {sorted(unmapped)}") | |
| 212 | return n | |
| 213 | ||
| 214 | ||
| 215 | def ingest_cbpd(conn, since): | |
| 216 | rows = [] | |
| 217 | for f in query_all(CBPD, since): | |
| 218 | a = f["attributes"] | |
| 219 | occurred = local_iso(a["cfs_datetime"]) | |
| 220 | if occurred is None: | |
| 221 | continue | |
| 222 | g = f.get("geometry") or {} | |
| 223 | code = a.get("incident_code") | |
| 224 | # Council Bluffs withholds the street address; the point is still exact. | |
| 225 | rows.append(("cbpd", a["cfs_number"], "Council Bluffs PD", | |
| 226 | a.get("case_number") or a.get("cfs_number"), occurred, | |
| 227 | a.get("incident_category"), code, a.get("disp_code"), None, | |
| 228 | int(code == CBPD_STOP_CODE), | |
| 229 | None, g.get("y"), g.get("x"))) | |
| 230 | return upsert(conn, rows) | |
| 231 | ||
| 232 | ||
| 233 | def ingest_opd_csv(conn, _since): | |
| 234 | """Backfill the 2015-2023 CSV archive. Statute text goes to offense_desc; | |
| 235 | category stays NULL because it is not a NIBRS category.""" | |
| 236 | rows = [] | |
| 237 | for path in sorted((ROOT / "raw_data").glob("Incidents_*.csv")): | |
| 238 | with path.open(newline="") as fh: | |
| 239 | if fh.readline().startswith("version https://git-lfs"): | |
| 240 | print(f" {path.name}: git-lfs pointer, run 'git lfs pull'") | |
| 241 | continue | |
| 242 | fh.seek(0) | |
| 243 | for i, r in enumerate(csv.reader(fh)): | |
| 244 | if len(r) < 8 or r[0] == "RB Number": | |
| 245 | continue | |
| 246 | rb, date, tm, desc, loc, district, lat, lon = r[:8] | |
| 247 | try: | |
| 248 | when = datetime.strptime(f"{date} {tm}", "%m/%d/%Y %H:%M") | |
| 249 | except ValueError: | |
| 250 | continue | |
| 251 | rows.append(("opd_csv", f"{path.stem}:{i}", "Omaha PD", rb, | |
| 252 | when.strftime("%Y-%m-%dT%H:%M:%S"), None, None, None, | |
| 253 | desc, 0, loc, | |
| 254 | float(lat) if lat else None, | |
| 255 | float(lon) if lon else None)) | |
| 256 | return upsert(conn, rows) | |
| 257 | ||
| 258 | ||
| 259 | def ingest_alpr(conn, _since): | |
| 260 | payload = get(OVERPASS, {"data": OVERPASS_QUERY}) | |
| 261 | now = datetime.now(LOCAL).strftime("%Y-%m-%dT%H:%M:%S") | |
| 262 | rows = [] | |
| 263 | for e in payload["elements"]: | |
| 264 | t = e.get("tags", {}) | |
| 265 | rows.append((e["id"], e["lat"], e["lon"], | |
| 266 | t.get("manufacturer") or t.get("brand"), | |
| 267 | t.get("operator"), | |
| 268 | t.get("direction") or t.get("camera:direction"), | |
| 269 | now, now, json.dumps(t, sort_keys=True))) | |
| 270 | conn.executemany( | |
| 271 | """INSERT INTO alpr_cameras | |
| 272 | (osm_id, lat, lon, manufacturer, operator, direction, | |
| 273 | first_seen, last_seen, tags) | |
| 274 | VALUES (?,?,?,?,?,?,?,?,?) | |
| 275 | ON CONFLICT(osm_id) DO UPDATE SET | |
| 276 | lat=excluded.lat, lon=excluded.lon, | |
| 277 | manufacturer=excluded.manufacturer, operator=excluded.operator, | |
| 278 | direction=excluded.direction, last_seen=excluded.last_seen, | |
| 279 | tags=excluded.tags""", | |
| 280 | rows) | |
| 281 | return len(rows) | |
| 282 | ||
| 283 | ||
| 284 | def migrate(conn): | |
| 285 | """Add columns introduced after a database was first built. Runs before | |
| 286 | schema.sql so its indexes can reference the new columns.""" | |
| 287 | have = {r[1] for r in conn.execute("PRAGMA table_info(incidents)")} | |
| 288 | if have and "is_stop" not in have: | |
| 289 | conn.execute("ALTER TABLE incidents ADD COLUMN is_stop INTEGER NOT NULL" | |
| 290 | " DEFAULT 0") | |
| 291 | conn.execute("UPDATE incidents SET is_stop = 1 WHERE category = ?", | |
| 292 | (SARPY_STOP_CATEGORY,)) | |
| 293 | conn.commit() | |
| 294 | ||
| 295 | ||
| 296 | SOURCES = { | |
| 297 | "opd": ingest_opd, | |
| 298 | "sarpy": ingest_sarpy, | |
| 299 | "cbpd": ingest_cbpd, | |
| 300 | "alpr": ingest_alpr, | |
| 301 | "opd_csv": ingest_opd_csv, | |
| 302 | } | |
| 303 | ||
| 304 | ||
| 305 | def main(): | |
| 306 | p = argparse.ArgumentParser(description=__doc__, | |
| 307 | formatter_class=argparse.RawDescriptionHelpFormatter) | |
| 308 | p.add_argument("sources", nargs="*", choices=list(SOURCES), | |
| 309 | help="default: opd sarpy cbpd alpr") | |
| 310 | p.add_argument("--since-days", type=int, default=30, | |
| 311 | help="only pull incidents this recent (default 30); " | |
| 312 | "ignored by alpr and opd_csv") | |
| 313 | p.add_argument("--full", action="store_true", | |
| 314 | help="pull the complete feed instead of --since-days") | |
| 315 | args = p.parse_args() | |
| 316 | sources = args.sources or ["opd", "sarpy", "cbpd", "alpr"] | |
| 317 | ||
| 318 | since = None if args.full else datetime.now(timezone.utc) - timedelta(days=args.since_days) | |
| 319 | ||
| 320 | DB.parent.mkdir(exist_ok=True) | |
| 321 | conn = sqlite3.connect(DB) | |
| 322 | migrate(conn) | |
| 323 | conn.executescript((ROOT / "schema.sql").read_text()) | |
| 324 | for name in sources: | |
| 325 | start = time.monotonic() | |
| 326 | print(f" {name}: pulling...") | |
| 327 | n = SOURCES[name](conn, since) | |
| 328 | conn.commit() | |
| 329 | print(f" {name}: {n} rows in {time.monotonic() - start:.1f}s") | |
| 330 | conn.close() | |
| 331 | ||
| 332 | ||
| 333 | if __name__ == "__main__": | |
| 334 | main() | |
requirements-ingest.txt added +1
| @@ -0,0 +1 @@ | ||
| 1 | certifi | |
requirements.txt added +5
| @@ -0,0 +1,5 @@ | ||
| 1 | -r requirements-ingest.txt | |
| 2 | dash>=3.0 | |
| 3 | numpy>=2.0 | |
| 4 | pandas>=2.2 | |
| 5 | plotly>=6.0 | |
schema.sql added +38
| @@ -0,0 +1,38 @@ | ||
| 1 | -- Incidents from every ingested feed. occurred_at is local time (America/Chicago); | |
| 2 | -- the upstream services return UTC epochs and ingest.py converts on the way in. | |
| 3 | CREATE TABLE IF NOT EXISTS incidents ( | |
| 4 | source TEXT NOT NULL, -- opd | sarpy | cbpd | opd_csv | |
| 5 | source_key TEXT NOT NULL, -- PK (opd) | IncidentId (sarpy) | cfs_number (cbpd) | |
| 6 | agency TEXT NOT NULL, | |
| 7 | case_id TEXT, | |
| 8 | occurred_at TEXT NOT NULL, -- ISO 8601, no offset | |
| 9 | category TEXT, -- each feed's own taxonomy, not comparable | |
| 10 | call_type TEXT, -- CadTypeDesc (sarpy) | incident_code (cbpd) | |
| 11 | disposition TEXT, -- CAD disposition, sarpy and cbpd | |
| 12 | offense_desc TEXT, -- StatuteDesc (sarpy) | statute text (opd_csv) | |
| 13 | is_stop INTEGER NOT NULL DEFAULT 0, -- officer-initiated vehicle stop | |
| 14 | address TEXT, | |
| 15 | lat REAL, | |
| 16 | lon REAL, | |
| 17 | PRIMARY KEY (source, source_key) | |
| 18 | ); | |
| 19 | ||
| 20 | CREATE INDEX IF NOT EXISTS incidents_occurred ON incidents (occurred_at); | |
| 21 | CREATE INDEX IF NOT EXISTS incidents_agency ON incidents (agency, occurred_at); | |
| 22 | CREATE INDEX IF NOT EXISTS incidents_category ON incidents (category); | |
| 23 | CREATE INDEX IF NOT EXISTS incidents_stop ON incidents (is_stop, occurred_at); | |
| 24 | ||
| 25 | -- ALPR cameras from OpenStreetMap (ODbL). first_seen/last_seen track when a node | |
| 26 | -- entered and was last present in the Overpass result, so cameras that appear or | |
| 27 | -- are removed are visible over time. | |
| 28 | CREATE TABLE IF NOT EXISTS alpr_cameras ( | |
| 29 | osm_id INTEGER PRIMARY KEY, | |
| 30 | lat REAL NOT NULL, | |
| 31 | lon REAL NOT NULL, | |
| 32 | manufacturer TEXT, | |
| 33 | operator TEXT, | |
| 34 | direction TEXT, | |
| 35 | first_seen TEXT NOT NULL, | |
| 36 | last_seen TEXT NOT NULL, | |
| 37 | tags TEXT NOT NULL -- full OSM tag dict as JSON | |
| 38 | ); | |