Skip to content

Latest commit

 

History

History
322 lines (263 loc) · 16 KB

File metadata and controls

322 lines (263 loc) · 16 KB

OnTime

The OnTime dataset has records of airline on-time performance of domestic US flights from 1987 to 2025. The total dataset has around 200 million rows.

Part 1: Getting data into MergeTree

From Parquet data in S3

Altinity provides ontime data as Parquet files available in a public S3 bucket at s3://altinity-playground-datasets/ontime/parts/*.parquet. This bucket is set up as requester pays, so you must authenticate with an AWS account to access it. To avoid this small cost, see how to import from raw data.

If you don't want the data in the default database, run CREATE DATABASE ontime and USE ontime, or use your desired name. Then, run the following SQL to create the table and insert the data using the S3 table function. Replace aws_access_key_id, aws_secret_access_key, and aws_session_token with your AWS credentials. If you don't have a session token, remove the parameter.

CREATE TABLE ontime
ENGINE = MergeTree
PARTITION BY Year
ORDER BY (Year, Quarter, Month, DayofMonth, FlightDate, IATA_CODE_Reporting_Airline)
AS
SELECT * FROM s3(
  's3://altinity-playground-datasets/ontime/parts/*.parquet',
  'aws_access_key_id', 'aws_secret_access_key', 'aws_session_token',
  'Parquet',
  'Year UInt16, Quarter UInt8, Month UInt8, DayofMonth UInt8, DayOfWeek UInt8, FlightDate Date, Reporting_Airline LowCardinality(String), DOT_ID_Reporting_Airline Int32, IATA_CODE_Reporting_Airline LowCardinality(String), Tail_Number LowCardinality(String), Flight_Number_Reporting_Airline LowCardinality(String), OriginAirportID Int32, OriginAirportSeqID Int32, OriginCityMarketID Int32, Origin FixedString(5), OriginCityName LowCardinality(String), OriginState FixedString(2), OriginStateFips FixedString(2), OriginStateName LowCardinality(String), OriginWac Int32, DestAirportID Int32, DestAirportSeqID Int32, DestCityMarketID Int32, Dest FixedString(5), DestCityName LowCardinality(String), DestState FixedString(2), DestStateFips FixedString(2), DestStateName LowCardinality(String), DestWac Int32, CRSDepTime Int32, DepTime Int32, DepDelay Int32, DepDelayMinutes Int32, DepDel15 Int32, DepartureDelayGroups LowCardinality(String), DepTimeBlk LowCardinality(String), TaxiOut Int32, WheelsOff LowCardinality(String), WheelsOn LowCardinality(String), TaxiIn Int32, CRSArrTime Int32, ArrTime Int32, ArrDelay Int32, ArrDelayMinutes Int32, ArrDel15 Int32, ArrivalDelayGroups LowCardinality(String), ArrTimeBlk LowCardinality(String), Cancelled Int8, CancellationCode FixedString(1), Diverted Int8, CRSElapsedTime Int32, ActualElapsedTime Int32, AirTime Int32, Flights Int32, Distance Int32, DistanceGroup Int8, CarrierDelay Int32, WeatherDelay Int32, NASDelay Int32, SecurityDelay Int32, LateAircraftDelay Int32, FirstDepTime Int16, TotalAddGTime Int16, LongestAddGTime Int16, DivAirportLandings Int8, DivReachedDest Int8, DivActualElapsedTime Int16, DivArrDelay Int16, DivDistance Int16, Div1Airport LowCardinality(String), Div1AirportID Int32, Div1AirportSeqID Int32, Div1WheelsOn Int16, Div1TotalGTime Int16, Div1LongestGTime Int16, Div1WheelsOff Int16, Div1TailNum LowCardinality(String), Div2Airport LowCardinality(String), Div2AirportID Int32, Div2AirportSeqID Int32, Div2WheelsOn Int16, Div2TotalGTime Int16, Div2LongestGTime Int16, Div2WheelsOff Int16, Div2TailNum LowCardinality(String), Div3Airport LowCardinality(String), Div3AirportID Int32, Div3AirportSeqID Int32, Div3WheelsOn Int16, Div3TotalGTime Int16, Div3LongestGTime Int16, Div3WheelsOff Int16, Div3TailNum LowCardinality(String), Div4Airport LowCardinality(String), Div4AirportID Int32, Div4AirportSeqID Int32, Div4WheelsOn Int16, Div4TotalGTime Int16, Div4LongestGTime Int16, Div4WheelsOff Int16, Div4TailNum LowCardinality(String), Div5Airport LowCardinality(String), Div5AirportID Int32, Div5AirportSeqID Int32, Div5WheelsOn Int16, Div5TotalGTime Int16, Div5LongestGTime Int16, Div5WheelsOff Int16, Div5TailNum LowCardinality(String)',
  headers('x-amz-request-payer' = 'requester')
)
SETTINGS max_threads=4, max_insert_threads=4, input_format_parallel_parsing=0;

From raw data

If you already have the ontime table in MergeTree loaded from S3, skip to Part 2. This part references the ClickHouse docs page. Otherwise instead of using the S3 bucket, you can download the data from the source:

wget --no-check-certificate --continue https://transtats.bts.gov/PREZIP/On_Time_Reporting_Carrier_On_Time_Performance_1987_present_{1987..2024}_{1..12}.zip
wget --no-check-certificate --continue https://transtats.bts.gov/PREZIP/On_Time_Reporting_Carrier_On_Time_Performance_1987_present_2025_{1..11}.zip

Data from 2025 is currently only available until November.

If the source data is unavailable, you can also download the raw data from Altinity's S3 bucket with requester pays (AWS account required):

aws s3 cp s3://altinity-playground-datasets/ontime/raw ./ --recursive --request-payer requester

Then, log into ClickHouse client and create the database and table below:

CREATE DATABASE ontime;

USE ontime;
CREATE TABLE `ontime`
(
    `Year`                            UInt16,
    `Quarter`                         UInt8,
    `Month`                           UInt8,
    `DayofMonth`                      UInt8,
    `DayOfWeek`                       UInt8,
    `FlightDate`                      Date,
    `Reporting_Airline`               LowCardinality(String),
    `DOT_ID_Reporting_Airline`        Int32,
    `IATA_CODE_Reporting_Airline`     LowCardinality(String),
    `Tail_Number`                     LowCardinality(String),
    `Flight_Number_Reporting_Airline` LowCardinality(String),
    `OriginAirportID`                 Int32,
    `OriginAirportSeqID`              Int32,
    `OriginCityMarketID`              Int32,
    `Origin`                          FixedString(5),
    `OriginCityName`                  LowCardinality(String),
    `OriginState`                     FixedString(2),
    `OriginStateFips`                 FixedString(2),
    `OriginStateName`                 LowCardinality(String),
    `OriginWac`                       Int32,
    `DestAirportID`                   Int32,
    `DestAirportSeqID`                Int32,
    `DestCityMarketID`                Int32,
    `Dest`                            FixedString(5),
    `DestCityName`                    LowCardinality(String),
    `DestState`                       FixedString(2),
    `DestStateFips`                   FixedString(2),
    `DestStateName`                   LowCardinality(String),
    `DestWac`                         Int32,
    `CRSDepTime`                      Int32,
    `DepTime`                         Int32,
    `DepDelay`                        Int32,
    `DepDelayMinutes`                 Int32,
    `DepDel15`                        Int32,
    `DepartureDelayGroups`            LowCardinality(String),
    `DepTimeBlk`                      LowCardinality(String),
    `TaxiOut`                         Int32,
    `WheelsOff`                       LowCardinality(String),
    `WheelsOn`                        LowCardinality(String),
    `TaxiIn`                          Int32,
    `CRSArrTime`                      Int32,
    `ArrTime`                         Int32,
    `ArrDelay`                        Int32,
    `ArrDelayMinutes`                 Int32,
    `ArrDel15`                        Int32,
    `ArrivalDelayGroups`              LowCardinality(String),
    `ArrTimeBlk`                      LowCardinality(String),
    `Cancelled`                       Int8,
    `CancellationCode`                FixedString(1),
    `Diverted`                        Int8,
    `CRSElapsedTime`                  Int32,
    `ActualElapsedTime`               Int32,
    `AirTime`                         Int32,
    `Flights`                         Int32,
    `Distance`                        Int32,
    `DistanceGroup`                   Int8,
    `CarrierDelay`                    Int32,
    `WeatherDelay`                    Int32,
    `NASDelay`                        Int32,
    `SecurityDelay`                   Int32,
    `LateAircraftDelay`               Int32,
    `FirstDepTime`                    Int16,
    `TotalAddGTime`                   Int16,
    `LongestAddGTime`                 Int16,
    `DivAirportLandings`              Int8,
    `DivReachedDest`                  Int8,
    `DivActualElapsedTime`            Int16,
    `DivArrDelay`                     Int16,
    `DivDistance`                     Int16,
    `Div1Airport`                     LowCardinality(String),
    `Div1AirportID`                   Int32,
    `Div1AirportSeqID`                Int32,
    `Div1WheelsOn`                    Int16,
    `Div1TotalGTime`                  Int16,
    `Div1LongestGTime`                Int16,
    `Div1WheelsOff`                   Int16,
    `Div1TailNum`                     LowCardinality(String),
    `Div2Airport`                     LowCardinality(String),
    `Div2AirportID`                   Int32,
    `Div2AirportSeqID`                Int32,
    `Div2WheelsOn`                    Int16,
    `Div2TotalGTime`                  Int16,
    `Div2LongestGTime`                Int16,
    `Div2WheelsOff`                   Int16,
    `Div2TailNum`                     LowCardinality(String),
    `Div3Airport`                     LowCardinality(String),
    `Div3AirportID`                   Int32,
    `Div3AirportSeqID`                Int32,
    `Div3WheelsOn`                    Int16,
    `Div3TotalGTime`                  Int16,
    `Div3LongestGTime`                Int16,
    `Div3WheelsOff`                   Int16,
    `Div3TailNum`                     LowCardinality(String),
    `Div4Airport`                     LowCardinality(String),
    `Div4AirportID`                   Int32,
    `Div4AirportSeqID`                Int32,
    `Div4WheelsOn`                    Int16,
    `Div4TotalGTime`                  Int16,
    `Div4LongestGTime`                Int16,
    `Div4WheelsOff`                   Int16,
    `Div4TailNum`                     LowCardinality(String),
    `Div5Airport`                     LowCardinality(String),
    `Div5AirportID`                   Int32,
    `Div5AirportSeqID`                Int32,
    `Div5WheelsOn`                    Int16,
    `Div5TotalGTime`                  Int16,
    `Div5LongestGTime`                Int16,
    `Div5WheelsOff`                   Int16,
    `Div5TailNum`                     LowCardinality(String)
) ENGINE = MergeTree
  PARTITION BY Year
  ORDER BY (Year, Quarter, Month, DayofMonth, FlightDate, IATA_CODE_Reporting_Airline);

Then, in the directory you have the zip files downloaded on your local machine, run:

ls -1 *.zip | xargs -I{} -P $(nproc) bash -c "echo {}; unzip -cq {} '*.csv' | sed 's/\.00//g' | clickhouse-client \
  --host your_hostname \
  --port 9000 \
  --user your_user \
  --password 'your_password' \
  --input_format_csv_empty_as_default 1 \
  --query='INSERT INTO ontime.ontime FORMAT CSVWithNames'"

Replace the values for --host, --port, --user, and --password as needed.

You should now see the rows in the ontime.ontime table:

SELECT COUNT(*) from ontime.ontime;

   ┌───count()─┐
1. │ 230307689 │ -- 230.31 million
   └───────────┘

1 row in set. Elapsed: 0.002 sec.

Part 2: Exporting data from MergeTree to Iceberg

Optional: Force merge parts

Before exporting the MergeTree parts, you can let ClickHouse merge parts so there is only 1 part per partition. This is done in the background and may take a while, so you can immediately merge all parts by running:

OPTIMIZE TABLE ontime.ontime FINAL;

These above commands will load your ClickHouse server heavily, so don't run them in production.

Generate commands for EXPORT PART

Using the script

You can use the export parts script to automate exporting parts from MergeTree to S3 in Parquet format. This script runs ALTER TABLE EXPORT PART queries in small batches, waiting for exports to complete before running new ones.

Manual commands

Otherwise, follow the steps below to manually run export part commands:

Replace s3://your-warehouse below with the S3 bucket for your environment's Iceberg "Warehouse" on the Catalogs page (where?).

SELECT 'ALTER TABLE ontime.ontime EXPORT PART \'' || name || '\' TO TABLE FUNCTION s3(\'s3://your-warehouse/ontime/ontime/data/Year={_partition_id}/{_file}.parquet\', format=\'parquet\', partition_strategy=\'wildcard\') PARTITION BY Year settings allow_experimental_export_merge_tree_part=1;'
FROM system.parts
WHERE active and (database, table) in ('ontime','ontime')
FORMAT TSVRaw;

Copy the first command and run it to make it sure it works before running the rest of them. Now, our Parquet files are in the location we will create the Iceberg table.

Create Iceberg table with ice

Once the data is in S3, we need to create the table in the Iceberg catalog.

Prerequisites:

  • The Ice REST catalog is set up for the environment in the ACM
  • The ice CLI is installed on your local machine

Create a config file .ice.yaml for the ice CLI tool. Copy this template and fill in the uri and bearerToken from the "Catalog URL" and "Auth Token" values from the Iceberg catalog connection details (where?).

uri: YOUR_CATALOG_URL
bearerToken: YOUR_AUTH_TOKEN

In the same directory as your .ice.yaml file, run the ice commands to register the tables to the catalog. Replace s3://your-warehouse with the "Warehouse" value on the Catalogs page.

ice insert ontime.ontime --insecure -p 's3://your-warehouse/ontime/ontime/data/*/*.parquet' --no-copy --thread-count=10 --partition='[{"column":"Year"}]'

Create the ice database in ClickHouse

If you don't have the ice database in ClickHouse already that points to your Iceberg data, follow the steps to create it and name the database "ice". Alternatively, you can run the following SQL commands in your ClickHouse cluster:

SET allow_experimental_database_iceberg = 1;

DROP DATABASE IF EXISTS ice;

CREATE DATABASE ice
  ENGINE = DataLakeCatalog('http://ice-rest-catalog:5000')
  SETTINGS catalog_type = 'rest',
    auth_header = 'Authorization: Bearer your-auth-token', 
    warehouse = 's3://your-warehouse';

Replace the values your-auth-token and s3://your-warehouse with the Catalog connection details.

Now, you should see your Iceberg table from ClickHouse:

SHOW TABLES IN ice;

SELECT count(*) from ice.`ontime.ontime`;

The ontime table is located in the ontime Iceberg namespace within the ice database, so you must use backticks.

Part 3: Create a Hybrid table

Hybrid tables allow you to query data with part of the data in MergeTree and part of it in Iceberg using a "watermark" condition. The main use case is to separate frequently accessed "hot" data in MergeTree from less-used "cold" data in Iceberg for cheaper object storage.

Note that we have the entire dataset in both MergeTree and Iceberg already for this demo. In a production environment you wouldn't have duplicate data in both.

The following SQL creates a table with the Hybrid engine and a specified watermark date. If your MergeTree table is not in a database called ontime, change it in the command. Replace the conditions Year >= 2015 and Year < 2015 to your desired watermark condition. You can reference any other column such as FlightDate.

CREATE TABLE ontime.ontime_hybrid AS ontime.ontime
ENGINE = Hybrid(
  cluster('{cluster}', ontime.ontime),
  Year >= 2015,
  ice.`ontime.ontime`,
  Year < 2015
)
SETTINGS allow_experimental_hybrid_table=1;

For more information on Hybrid tables, see this blog post.