-
Notifications
You must be signed in to change notification settings - Fork 2
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
FEAT: Add 'revenue' field to Rail data pipeline (#446)
This change adds the GTFS-RT "revenue" field to the LAMP Rail Performance Data pipeline. - Create db migrations to add revenue column to vehicle_trips table - Read vehicle.trip.revenue column from VehiclePositions parquet files - Load revenue data from parquet files into DB table - Add is_revenue column to LAMP_ALL_RT_fields OPMI export - Add where clause for only revenue trips to Subway Performance Data export The use of the revenue field differs between the OPMI export and Subway Performance Data export. For the OMPI export an additional field is added so that revenue trips can be filtered out at their discretion. The Subway Performance Data export was never meant to include non-revenue event data, so it is filtered out from the export all together. Asana Task: https://app.asana.com/0/1205827492903547/1208216161546522
- Loading branch information
Showing
19 changed files
with
306 additions
and
141 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
74 changes: 74 additions & 0 deletions
74
...mp_py/migrations/versions/performance_manager_dev/008_32ba735d080c_add_revenue_columns.py
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,74 @@ | ||
"""add revenue columns | ||
Revision ID: 32ba735d080c | ||
Revises: 896dedd8a4db | ||
Create Date: 2024-09-20 08:47:52.784591 | ||
This change adds a boolean revenue column to the vehcile_trips table. | ||
Initially this will be filled with True and back-filled by a seperate operation | ||
Details | ||
* upgrade -> drop triggers and indexes from table and add revenue column | ||
* downgrade -> drop revenue column | ||
""" | ||
|
||
from alembic import op | ||
import sqlalchemy as sa | ||
|
||
from lamp_py.postgres.rail_performance_manager_schema import ( | ||
TempEventCompare, | ||
VehicleTrips, | ||
) | ||
|
||
# revision identifiers, used by Alembic. | ||
revision = "32ba735d080c" | ||
down_revision = "896dedd8a4db" | ||
branch_labels = None | ||
depends_on = None | ||
|
||
|
||
def upgrade() -> None: | ||
op.execute( | ||
f"ALTER TABLE public.vehicle_trips DISABLE TRIGGER rt_trips_update_branch_trunk;" | ||
) | ||
op.execute( | ||
f"ALTER TABLE public.vehicle_trips DISABLE TRIGGER update_vehicle_trips_modified;" | ||
) | ||
op.drop_index("ix_vehicle_trips_composite_1", table_name="vehicle_trips") | ||
op.drop_constraint("vehicle_trips_unique_trip", table_name="vehicle_trips") | ||
|
||
op.add_column( | ||
"temp_event_compare", sa.Column("revenue", sa.Boolean(), nullable=True) | ||
) | ||
op.add_column( | ||
"vehicle_trips", sa.Column("revenue", sa.Boolean(), nullable=True) | ||
) | ||
op.execute(sa.update(TempEventCompare).values(revenue=True)) | ||
op.execute(sa.update(VehicleTrips).values(revenue=True)) | ||
op.alter_column("temp_event_compare", "revenue", nullable=False) | ||
op.alter_column("vehicle_trips", "revenue", nullable=False) | ||
|
||
op.create_unique_constraint( | ||
"vehicle_trips_unique_trip", | ||
"vehicle_trips", | ||
["service_date", "route_id", "trip_id"], | ||
) | ||
op.create_index( | ||
"ix_vehicle_trips_composite_1", | ||
"vehicle_trips", | ||
["route_id", "direction_id", "vehicle_id"], | ||
unique=False, | ||
) | ||
op.execute( | ||
f"ALTER TABLE public.vehicle_trips ENABLE TRIGGER rt_trips_update_branch_trunk;" | ||
) | ||
op.execute( | ||
f"ALTER TABLE public.vehicle_trips ENABLE TRIGGER update_vehicle_trips_modified;" | ||
) | ||
|
||
|
||
def downgrade() -> None: | ||
op.drop_column("vehicle_trips", "revenue") | ||
op.drop_column("temp_event_compare", "revenue") |
Oops, something went wrong.