20 KiB
Postgres Integration
This document is the canonical table contract for the optional postgres sink.
It describes the schema created and written by weatherfeeder through feedkit's
Postgres sink.
Configure the sink as described in configuration.
Initialization And Writes
At startup, each configured Postgres sink opens the database and runs
CREATE TABLE IF NOT EXISTS for every weatherfeeder table, followed by
CREATE INDEX IF NOT EXISTS for every configured index.
This initialization creates missing tables and indexes only. It does not alter existing tables, migrate column definitions, drop old objects, or backfill data. Schema changes require operator-managed database migration.
Events are mapped only for canonical weather schemas:
weather.observation.v1weather.forecast.v1weather.forecast_discussion.v1weather.weather_story.v1weather.alert.v1weather.outlook.v1
Unsupported schemas produce no writes for this sink. Mapped events are inserted
transactionally. Inserts use ordinary INSERT; duplicate primary keys fail the
write.
Shared Envelope Columns
Parent tables store the feed event envelope:
| Column | Type | Null | Source |
|---|---|---|---|
event_id |
TEXT |
no | event.id |
event_kind |
TEXT |
no | event.kind |
event_source |
TEXT |
no | event.source |
event_schema |
TEXT |
no | event.schema |
event_emitted_at |
TIMESTAMPTZ |
no | event.emitted_at |
event_effective_at |
TIMESTAMPTZ |
yes | event.effective_at |
Table Overview
| Table | Primary key | Prune column |
|---|---|---|
observations |
event_id |
observed_at |
observation_present_weather |
event_id, weather_index |
observed_at |
forecasts |
event_id |
issued_at |
forecast_periods |
run_event_id, period_index |
issued_at |
forecast_discussions |
event_id |
issued_at |
forecast_discussion_key_messages |
run_event_id, message_index |
issued_at |
weather_story_runs |
event_id |
as_of |
weather_stories |
run_event_id, story_index |
as_of |
alert_runs |
event_id |
as_of |
alerts |
run_event_id, alert_index |
as_of |
alert_references |
run_event_id, alert_index, reference_index |
as_of |
outlook_runs |
event_id |
as_of |
outlooks |
run_event_id, outlook_index |
as_of |
Table Contract
observations
Primary key: event_id
Prune column: observed_at
Indexes:
idx_wf_obs_station_observed_atonstation_id,observed_atidx_wf_obs_observed_atonobserved_atidx_wf_obs_condition_codeoncondition_code
| Column | Type | Null | Source |
|---|---|---|---|
event_id |
TEXT |
no | event.id |
event_kind |
TEXT |
no | event.kind |
event_source |
TEXT |
no | event.source |
event_schema |
TEXT |
no | event.schema |
event_emitted_at |
TIMESTAMPTZ |
no | event.emitted_at |
event_effective_at |
TIMESTAMPTZ |
yes | event.effective_at |
station_id |
TEXT |
yes | payload.stationId |
station_name |
TEXT |
yes | payload.stationName |
observed_at |
TIMESTAMPTZ |
no | payload.timestamp |
condition_code |
INTEGER |
no | payload.conditionCode |
is_day |
BOOLEAN |
yes | payload.isDay |
text_description |
TEXT |
yes | payload.textDescription |
temperature_c |
DOUBLE PRECISION |
yes | payload.temperatureC |
dewpoint_c |
DOUBLE PRECISION |
yes | payload.dewpointC |
wind_direction_degrees |
DOUBLE PRECISION |
yes | payload.windDirectionDegrees |
wind_speed_kmh |
DOUBLE PRECISION |
yes | payload.windSpeedKmh |
wind_gust_kmh |
DOUBLE PRECISION |
yes | payload.windGustKmh |
barometric_pressure_pa |
DOUBLE PRECISION |
yes | payload.barometricPressurePa |
visibility_meters |
DOUBLE PRECISION |
yes | payload.visibilityMeters |
relative_humidity_percent |
DOUBLE PRECISION |
yes | payload.relativeHumidityPercent |
apparent_temperature_c |
DOUBLE PRECISION |
yes | payload.apparentTemperatureC |
observation_present_weather
Primary key: event_id, weather_index
Prune column: observed_at
Foreign key: event_id references observations(event_id) with cascade delete.
Index: idx_wf_obs_present_observed_at on observed_at
| Column | Type | Null | Source |
|---|---|---|---|
event_id |
TEXT REFERENCES observations(event_id) ON DELETE CASCADE |
no | Parent event ID. |
weather_index |
INTEGER |
no | payload.presentWeather[] index. |
observed_at |
TIMESTAMPTZ |
no | payload.timestamp |
raw_text |
TEXT |
yes | Compact JSON text from payload.presentWeather[].raw |
forecasts
Primary key: event_id
Prune column: issued_at
Indexes:
idx_wf_fc_location_product_issued_atonlocation_id,product,issued_atidx_wf_fc_issued_atonissued_atidx_wf_fc_product_issued_atonproduct,issued_at
| Column | Type | Null | Source |
|---|---|---|---|
event_id |
TEXT |
no | event.id |
event_kind |
TEXT |
no | event.kind |
event_source |
TEXT |
no | event.source |
event_schema |
TEXT |
no | event.schema |
event_emitted_at |
TIMESTAMPTZ |
no | event.emitted_at |
event_effective_at |
TIMESTAMPTZ |
yes | event.effective_at |
location_id |
TEXT |
yes | payload.locationId |
location_name |
TEXT |
yes | payload.locationName |
issued_at |
TIMESTAMPTZ |
no | payload.issuedAt |
updated_at |
TIMESTAMPTZ |
yes | payload.updatedAt |
product |
TEXT |
no | payload.product |
latitude |
DOUBLE PRECISION |
yes | payload.latitude |
longitude |
DOUBLE PRECISION |
yes | payload.longitude |
elevation_meters |
DOUBLE PRECISION |
yes | payload.elevationMeters |
period_count |
INTEGER |
no | len(payload.periods) |
forecast_periods
Primary key: run_event_id, period_index
Prune column: issued_at
Foreign key: run_event_id references forecasts(event_id) with cascade delete.
Indexes:
idx_wf_fc_period_start_timeonstart_timeidx_wf_fc_period_end_timeonend_timeidx_wf_fc_period_run_startonrun_event_id,start_time
| Column | Type | Null | Source |
|---|---|---|---|
run_event_id |
TEXT REFERENCES forecasts(event_id) ON DELETE CASCADE |
no | Parent event ID. |
period_index |
INTEGER |
no | payload.periods[] index. |
issued_at |
TIMESTAMPTZ |
no | Parent payload.issuedAt |
start_time |
TIMESTAMPTZ |
no | payload.periods[].startTime |
end_time |
TIMESTAMPTZ |
no | payload.periods[].endTime |
name |
TEXT |
yes | payload.periods[].name |
is_day |
BOOLEAN |
yes | payload.periods[].isDay |
condition_code |
INTEGER |
yes | payload.periods[].conditionCode |
text_description |
TEXT |
yes | payload.periods[].textDescription |
temperature_c |
DOUBLE PRECISION |
yes | payload.periods[].temperatureC |
temperature_c_min |
DOUBLE PRECISION |
yes | payload.periods[].temperatureCMin |
temperature_c_max |
DOUBLE PRECISION |
yes | payload.periods[].temperatureCMax |
dewpoint_c |
DOUBLE PRECISION |
yes | payload.periods[].dewpointC |
relative_humidity_percent |
DOUBLE PRECISION |
yes | payload.periods[].relativeHumidityPercent |
wind_direction_degrees |
DOUBLE PRECISION |
yes | payload.periods[].windDirectionDegrees |
wind_speed_kmh |
DOUBLE PRECISION |
yes | payload.periods[].windSpeedKmh |
wind_gust_kmh |
DOUBLE PRECISION |
yes | payload.periods[].windGustKmh |
barometric_pressure_pa |
DOUBLE PRECISION |
yes | payload.periods[].barometricPressurePa |
visibility_meters |
DOUBLE PRECISION |
yes | payload.periods[].visibilityMeters |
apparent_temperature_c |
DOUBLE PRECISION |
yes | payload.periods[].apparentTemperatureC |
cloud_cover_percent |
DOUBLE PRECISION |
yes | payload.periods[].cloudCoverPercent |
probability_of_precipitation_percent |
DOUBLE PRECISION |
yes | payload.periods[].probabilityOfPrecipitationPercent |
precipitation_amount_mm |
DOUBLE PRECISION |
yes | payload.periods[].precipitationAmountMm |
snowfall_depth_mm |
DOUBLE PRECISION |
yes | payload.periods[].snowfallDepthMm |
uv_index |
DOUBLE PRECISION |
yes | payload.periods[].uvIndex |
forecast_discussions
Primary key: event_id
Prune column: issued_at
Indexes:
idx_wf_discussion_office_product_issued_atonoffice_id,product,issued_atidx_wf_discussion_issued_atonissued_at
| Column | Type | Null | Source |
|---|---|---|---|
event_id |
TEXT |
no | event.id |
event_kind |
TEXT |
no | event.kind |
event_source |
TEXT |
no | event.source |
event_schema |
TEXT |
no | event.schema |
event_emitted_at |
TIMESTAMPTZ |
no | event.emitted_at |
event_effective_at |
TIMESTAMPTZ |
yes | event.effective_at |
office_id |
TEXT |
yes | payload.officeId |
office_name |
TEXT |
yes | payload.officeName |
issued_at |
TIMESTAMPTZ |
no | payload.issuedAt |
updated_at |
TIMESTAMPTZ |
yes | payload.updatedAt |
product |
TEXT |
no | payload.product |
short_term_qualifier |
TEXT |
yes | payload.shortTerm.qualifier |
short_term_issued_at |
TIMESTAMPTZ |
yes | payload.shortTerm.issuedAt |
short_term_text |
TEXT |
yes | payload.shortTerm.text |
long_term_qualifier |
TEXT |
yes | payload.longTerm.qualifier |
long_term_issued_at |
TIMESTAMPTZ |
yes | payload.longTerm.issuedAt |
long_term_text |
TEXT |
yes | payload.longTerm.text |
key_message_count |
INTEGER |
no | len(payload.keyMessages) |
forecast_discussion_key_messages
Primary key: run_event_id, message_index
Prune column: issued_at
Foreign key: run_event_id references forecast_discussions(event_id) with
cascade delete.
Index: idx_wf_discussion_message_issued_at on issued_at
| Column | Type | Null | Source |
|---|---|---|---|
run_event_id |
TEXT REFERENCES forecast_discussions(event_id) ON DELETE CASCADE |
no | Parent event ID. |
message_index |
INTEGER |
no | payload.keyMessages[] index. |
issued_at |
TIMESTAMPTZ |
no | Parent payload.issuedAt |
message_text |
TEXT |
yes | payload.keyMessages[] value |
weather_story_runs
Primary key: event_id
Prune column: as_of
Indexes:
idx_wf_story_run_office_as_ofonoffice_id,as_ofidx_wf_story_run_as_ofonas_of
| Column | Type | Null | Source |
|---|---|---|---|
event_id |
TEXT |
no | event.id |
event_kind |
TEXT |
no | event.kind |
event_source |
TEXT |
no | event.source |
event_schema |
TEXT |
no | event.schema |
event_emitted_at |
TIMESTAMPTZ |
no | event.emitted_at |
event_effective_at |
TIMESTAMPTZ |
yes | event.effective_at |
office_id |
TEXT |
yes | payload.officeId |
as_of |
TIMESTAMPTZ |
no | payload.asOf |
story_count |
INTEGER |
no | len(payload.stories) |
weather_stories
Primary key: run_event_id, story_index
Prune column: as_of
Foreign key: run_event_id references weather_story_runs(event_id) with
cascade delete.
Indexes:
idx_wf_stories_start_timeonstart_timeidx_wf_stories_end_timeonend_timeidx_wf_stories_updated_atonupdated_at
| Column | Type | Null | Source |
|---|---|---|---|
run_event_id |
TEXT REFERENCES weather_story_runs(event_id) ON DELETE CASCADE |
no | Parent event ID. |
story_index |
INTEGER |
no | payload.stories[] index. |
as_of |
TIMESTAMPTZ |
no | Parent payload.asOf |
office_id |
TEXT |
yes | payload.stories[].officeId |
start_time |
TIMESTAMPTZ |
no | payload.stories[].startTime |
end_time |
TIMESTAMPTZ |
no | payload.stories[].endTime |
updated_at |
TIMESTAMPTZ |
no | payload.stories[].updatedAt |
title |
TEXT |
yes | payload.stories[].title |
description |
TEXT |
yes | payload.stories[].description |
alt_text |
TEXT |
yes | payload.stories[].altText |
priority |
BOOLEAN |
no | payload.stories[].priority |
story_order |
INTEGER |
no | payload.stories[].order |
download_url |
TEXT |
yes | payload.stories[].downloadUrl |
alert_runs
Primary key: event_id
Prune column: as_of
Indexes:
idx_wf_alert_run_location_as_ofonlocation_id,as_ofidx_wf_alert_run_as_ofonas_of
| Column | Type | Null | Source |
|---|---|---|---|
event_id |
TEXT |
no | event.id |
event_kind |
TEXT |
no | event.kind |
event_source |
TEXT |
no | event.source |
event_schema |
TEXT |
no | event.schema |
event_emitted_at |
TIMESTAMPTZ |
no | event.emitted_at |
event_effective_at |
TIMESTAMPTZ |
yes | event.effective_at |
location_id |
TEXT |
yes | payload.locationId |
location_name |
TEXT |
yes | payload.locationName |
as_of |
TIMESTAMPTZ |
no | payload.asOf |
latitude |
DOUBLE PRECISION |
yes | payload.latitude |
longitude |
DOUBLE PRECISION |
yes | payload.longitude |
alert_count |
INTEGER |
no | len(payload.alerts) |
alerts
Primary key: run_event_id, alert_index
Prune column: as_of
Foreign key: run_event_id references alert_runs(event_id) with cascade
delete.
Indexes:
idx_wf_alerts_alert_idonalert_ididx_wf_alerts_severity_expiresonseverity,expiresidx_wf_alerts_as_ofonas_of
| Column | Type | Null | Source |
|---|---|---|---|
run_event_id |
TEXT REFERENCES alert_runs(event_id) ON DELETE CASCADE |
no | Parent event ID. |
alert_index |
INTEGER |
no | payload.alerts[] index. |
as_of |
TIMESTAMPTZ |
no | Parent payload.asOf |
alert_id |
TEXT |
no | payload.alerts[].id |
event |
TEXT |
yes | payload.alerts[].event |
headline |
TEXT |
yes | payload.alerts[].headline |
severity |
TEXT |
yes | payload.alerts[].severity |
urgency |
TEXT |
yes | payload.alerts[].urgency |
certainty |
TEXT |
yes | payload.alerts[].certainty |
status |
TEXT |
yes | payload.alerts[].status |
message_type |
TEXT |
yes | payload.alerts[].messageType |
category |
TEXT |
yes | payload.alerts[].category |
response |
TEXT |
yes | payload.alerts[].response |
description |
TEXT |
yes | payload.alerts[].description |
instruction |
TEXT |
yes | payload.alerts[].instruction |
sent |
TIMESTAMPTZ |
yes | payload.alerts[].sent |
effective |
TIMESTAMPTZ |
yes | payload.alerts[].effective |
onset |
TIMESTAMPTZ |
yes | payload.alerts[].onset |
expires |
TIMESTAMPTZ |
yes | payload.alerts[].expires |
area_description |
TEXT |
yes | payload.alerts[].areaDescription |
sender_name |
TEXT |
yes | payload.alerts[].senderName |
reference_count |
INTEGER |
no | len(payload.alerts[].references) |
alert_references
Primary key: run_event_id, alert_index, reference_index
Prune column: as_of
Foreign key: run_event_id references alert_runs(event_id) with cascade
delete.
Indexes:
idx_wf_alert_refs_as_ofonas_ofidx_wf_alert_refs_sentonsent
| Column | Type | Null | Source |
|---|---|---|---|
run_event_id |
TEXT REFERENCES alert_runs(event_id) ON DELETE CASCADE |
no | Parent event ID. |
alert_index |
INTEGER |
no | Parent alert index. |
reference_index |
INTEGER |
no | payload.alerts[].references[] index. |
as_of |
TIMESTAMPTZ |
no | Parent payload.asOf |
id |
TEXT |
yes | payload.alerts[].references[].id |
identifier |
TEXT |
yes | payload.alerts[].references[].identifier |
sender |
TEXT |
yes | payload.alerts[].references[].sender |
sent |
TIMESTAMPTZ |
yes | payload.alerts[].references[].sent |
outlook_runs
Primary key: event_id
Prune column: as_of
Indexes:
idx_wf_outlook_run_location_as_ofonlocation_id,as_ofidx_wf_outlook_run_as_ofonas_of
| Column | Type | Null | Source |
|---|---|---|---|
event_id |
TEXT |
no | event.id |
event_kind |
TEXT |
no | event.kind |
event_source |
TEXT |
no | event.source |
event_schema |
TEXT |
no | event.schema |
event_emitted_at |
TIMESTAMPTZ |
no | event.emitted_at |
event_effective_at |
TIMESTAMPTZ |
yes | event.effective_at |
location_id |
TEXT |
yes | payload.locationId |
location_name |
TEXT |
yes | payload.locationName |
latitude |
DOUBLE PRECISION |
yes | payload.latitude |
longitude |
DOUBLE PRECISION |
yes | payload.longitude |
as_of |
TIMESTAMPTZ |
no | payload.asOf |
issued_at |
TIMESTAMPTZ |
yes | payload.issuedAt |
outlook_count |
INTEGER |
no | len(payload.outlooks) |
outlooks
Primary key: run_event_id, outlook_index
Prune column: as_of
Foreign key: run_event_id references outlook_runs(event_id) with cascade
delete.
Indexes:
idx_wf_outlooks_contains_validoncontains_location,valid_from,valid_toidx_wf_outlooks_day_type_labelonday,outlook_type,labelidx_wf_outlooks_validonvalid_from,valid_to
| Column | Type | Null | Source |
|---|---|---|---|
run_event_id |
TEXT REFERENCES outlook_runs(event_id) ON DELETE CASCADE |
no | Parent event ID. |
outlook_index |
INTEGER |
no | payload.outlooks[] index. |
as_of |
TIMESTAMPTZ |
no | Parent payload.asOf |
product |
TEXT |
no | payload.outlooks[].product |
day |
INTEGER |
no | payload.outlooks[].day |
outlook_type |
TEXT |
no | payload.outlooks[].outlookType |
label |
TEXT |
no | payload.outlooks[].label |
label_text |
TEXT |
yes | payload.outlooks[].labelText |
severity_rank |
INTEGER |
yes | payload.outlooks[].severityRank |
valid_from |
TIMESTAMPTZ |
no | payload.outlooks[].validFrom |
valid_to |
TIMESTAMPTZ |
no | payload.outlooks[].validTo |
issued_at |
TIMESTAMPTZ |
no | payload.outlooks[].issuedAt |
expires_at |
TIMESTAMPTZ |
no | payload.outlooks[].expiresAt |
forecaster |
TEXT |
yes | payload.outlooks[].forecaster |
headline |
TEXT |
yes | payload.outlooks[].headline |
summary |
TEXT |
yes | payload.outlooks[].summary |
discussion |
TEXT |
yes | payload.outlooks[].discussion |
source_url |
TEXT |
yes | payload.outlooks[].sourceUrl |
image_url |
TEXT |
yes | payload.outlooks[].imageUrl |
contains_location |
BOOLEAN |
no | payload.outlooks[].containsLocation |
geometry_json |
TEXT |
no | Compact JSON from payload.outlooks[].geometry |
Retention
When sink param prune is set, every successful write transaction deletes rows
older than now - prune from every table using that table's prune column.
The sink also exposes manual prune helpers in code, but the weatherfeeder
binary does not provide CLI commands for them.
Reconstructing Canonical Payloads
WeatherObservation: readobservations, then joinobservation_present_weatherbyevent_idordered byweather_index.WeatherForecastRun: readforecasts, then joinforecast_periodsbyrun_event_idordered byperiod_index.WeatherForecastDiscussion: readforecast_discussions, then joinforecast_discussion_key_messagesbyrun_event_idordered bymessage_index.WeatherStoryRun: readweather_story_runs, then joinweather_storiesbyrun_event_idordered bystory_index.WeatherAlertRun: readalert_runs, joinalertsbyrun_event_idordered byalert_index, then joinalert_referencesbyrun_event_idandalert_indexordered byreference_index.WeatherOutlookRun: readoutlook_runs, then joinoutlooksbyrun_event_idordered byoutlook_index.