diff --git a/prototypes/01_dataset_transform.py b/prototypes/01-1_dataset_transform.py similarity index 100% rename from prototypes/01_dataset_transform.py rename to prototypes/01-1_dataset_transform.py diff --git a/prototypes/01-2_dataset_transform_whole.py b/prototypes/01-2_dataset_transform_whole.py new file mode 100644 index 0000000..432c0ef --- /dev/null +++ b/prototypes/01-2_dataset_transform_whole.py @@ -0,0 +1,168 @@ +# %% +import datetime +import json +import pprint +import re +from collections import Counter +from pathlib import Path + +import polars as pl + +# %% +p_data_base = (Path.cwd() / "../data/Datenauszug_20260128").resolve() +assert p_data_base.exists() + +print("Total number of JSON files") +len(tuple(p_data_base.glob("**/*.json"))) + +# %% +# // check size contents of folder +folders = tuple(p_data_base.glob("*")) +folder_items = {} + +for folder in folders: + num_elements = len(tuple(folder.glob("*.json"))) + if num_elements == 0: + continue + folder_items[folder] = num_elements + print(f"Folder {folder} contains: {num_elements} items") +# %% +max_item_folder = Path() +max_val = float("-inf") + +for k, v in folder_items.items(): + if v > max_val: + max_item_folder = k + max_val = v + +print(f"Maximum number of items : {max_val}\nFolder: {max_item_folder}") +# %% +pprint.pprint(folder_items) + +# %% +# analyse type numbers for each folder +folder_to_types = [] +for idx, folder in enumerate(folders): + folder_types = [] + for file in folder.glob("*.json"): + with open(file, "r") as f: + data = json.load(f) + + type_num = data["initial"]["dsc_TypeNumber"]["value"] + folder_types.append(type_num) + + type_num_count = Counter(folder_types) + folder_to_types.append((folder.name, type_num_count)) + +typenum_counter = Counter() +for idx in range(len(folder_to_types)): + typenum_counter.update(folder_to_types[idx][1]) + +max_type_num, type_num_count = typenum_counter.most_common(1)[0] +print(f"Max type number is: {max_type_num}") +print(f"Number of occurrences: {type_num_count}") +# %% +typenum_counter +# %% +folder_to_types +# %% +# ** one file is one curve +# concatenate all files in one table +# %% +# one time-series entry is defined by its schema: +# [ts, ps, pressure, valve] +# [timestep, process_step, pressure_value, valve_value] +# valid states are ps = [101, 102, 110] + +BREAK_EARLY = False + +schema_read = { + "DU1210": pl.Float64, + "DU1260": pl.Float64, + "DU1270": pl.Float64, + "PRV1270": pl.Int64, + "V1210": pl.Boolean, + "V1250": pl.Boolean, + "V1264": pl.Boolean, + "V1270": pl.Boolean, + "V1510": pl.Boolean, + "V1560": pl.Boolean, + "VZ1210": pl.Float64, + "ps": pl.UInt32, + "ts": pl.String, + "type_num": pl.UInt8, + "id": pl.UInt64, +} +schema = schema_read.copy() +schema.update( + { + "ts": pl.Datetime, + "id": pl.UInt64, + "ts_delta_step": pl.Duration, + "ts_delta_cum": pl.Duration, + } +) +print(schema) +df = pl.DataFrame(schema=schema).with_columns(pl.col("ts").dt.replace_time_zone("UTC")) +count = 0 +for idx, file in enumerate(p_data_base.glob("**/*.json"), start=1): + with open(file, "r") as f: + data = json.load(f) + + type_num = data["initial"]["dsc_TypeNumber"]["value"] + df_file = pl.DataFrame(data["rows"], schema_overrides=schema_read) + df_file = df_file.with_columns( + pl.col("ts").str.to_datetime(time_zone="UTC"), + pl.lit(type_num).alias("type_num").cast(pl.UInt8), + pl.lit(idx).alias("id").cast(pl.UInt64), + ) + df_file = df_file.with_columns( + (pl.col.ts - pl.col.ts.shift(1)) + .alias("ts_delta_step") + .fill_null(pl.lit(0).cast(pl.Duration)) + ) + df_file = df_file.with_columns( + pl.col("ts_delta_step").cum_sum().alias("ts_delta_cum"), + ) + df = pl.concat((df, df_file)) + count += 1 + + if BREAK_EARLY and idx == 3: + break + +# df = df.with_columns(pl.col("ts").str.to_datetime(time_zone="UTC")) +# df = df.select( +# ["id", "type_num", "ts", "ts_delta_step", "ts_delta_cum", "ps", "DU1260", "V1560"] +# ) +df.head() +# %% +t = df.columns[0] + +sensor_feats: list[str] = [] +pattern = re.compile(r"^([A-Z]+[0-9]+)") +for feat in df.columns: + matches = pattern.match(feat) + if matches is None: + continue + + sensor_feats.append(matches.group(1)) + +sensor_feats.sort() + +feats_sorted = ["id", "type_num", "ts", "ts_delta_step", "ts_delta_cum", "ps"] + sensor_feats + +df = df.select(feats_sorted) +df.head() +# %% +print(f"Files processed: {count}") +print(f"Length of obtained data: {df.height}") +# %% +WRITE_TO_DISK = False + +concat_data = p_data_base / "all_data.parquet" +if WRITE_TO_DISK: + df.write_parquet(concat_data) +else: + df = pl.read_parquet(concat_data) + +# %%