Refactor CSV to Parquet conversion logic
This commit is contained in:
@@ -36,10 +36,8 @@ def run_pipeline():
|
|||||||
parse_log(RAW_LOG, BUS1_CSV, BUS2_CSV)
|
parse_log(RAW_LOG, BUS1_CSV, BUS2_CSV)
|
||||||
|
|
||||||
print("Converting to parquet...")
|
print("Converting to parquet...")
|
||||||
lf1 = parse_csv(BUS1_CSV)
|
parse_csv(BUS1_CSV).sink_parquet(BUS1_PARQUET)
|
||||||
lf1.sink_parquet(BUS1_PARQUET)
|
parse_csv(BUS2_CSV).sink_parquet(BUS2_PARQUET)
|
||||||
lf2 = parse_csv(BUS2_CSV)
|
|
||||||
lf2.sink_parquet(BUS2_PARQUET)
|
|
||||||
|
|
||||||
print("Decoding J1939...")
|
print("Decoding J1939...")
|
||||||
df1 = pl.read_parquet(BUS1_PARQUET)
|
df1 = pl.read_parquet(BUS1_PARQUET)
|
||||||
|
|||||||
Reference in New Issue
Block a user