Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 9a5a231841 | |||
| 74b393934c | |||
| 4a701e34c5 |
+32
-25
@@ -6,11 +6,7 @@ def get_j1939_mask() -> pl.Expr:
|
|||||||
Returns a Polars expression representing the strict J1939 filtering rules.
|
Returns a Polars expression representing the strict J1939 filtering rules.
|
||||||
"""
|
"""
|
||||||
id_int = pl.col("ID").str.to_integer(base=16).cast(pl.UInt32)
|
id_int = pl.col("ID").str.to_integer(base=16).cast(pl.UInt32)
|
||||||
return (
|
return id_int > 0x7FF
|
||||||
(id_int > 0x7FF) &
|
|
||||||
((id_int % 33554432 // 16777216) == 0) &
|
|
||||||
(pl.col("DLC") <= 8)
|
|
||||||
)
|
|
||||||
|
|
||||||
def decode_j1939_metadata(lf: pl.LazyFrame) -> pl.LazyFrame:
|
def decode_j1939_metadata(lf: pl.LazyFrame) -> pl.LazyFrame:
|
||||||
"""
|
"""
|
||||||
@@ -18,17 +14,18 @@ def decode_j1939_metadata(lf: pl.LazyFrame) -> pl.LazyFrame:
|
|||||||
"""
|
"""
|
||||||
id_int = pl.col("ID").str.to_integer(base=16).cast(pl.UInt32)
|
id_int = pl.col("ID").str.to_integer(base=16).cast(pl.UInt32)
|
||||||
|
|
||||||
id_shifted_8 = id_int // 256
|
|
||||||
id_shifted_16 = id_int // 65536
|
|
||||||
|
|
||||||
priority = ((id_int // 67108864) % 8).cast(pl.UInt8)
|
priority = ((id_int // 67108864) % 8).cast(pl.UInt8)
|
||||||
pf = (id_shifted_16 % 256).cast(pl.UInt8)
|
pf = ((id_int // 65536) % 256).cast(pl.UInt8)
|
||||||
ps = (id_shifted_8 % 256).cast(pl.UInt8)
|
ps = ((id_int // 256) % 256).cast(pl.UInt8)
|
||||||
sa = (id_int % 256).cast(pl.UInt8)
|
sa = (id_int % 256).cast(pl.UInt8)
|
||||||
|
|
||||||
da = pl.when(pf < 240).then(ps).otherwise(pl.lit(255, dtype=pl.UInt8))
|
da = pl.when(pf < 240).then(ps).otherwise(pl.lit(255, dtype=pl.UInt8)).cast(pl.UInt8)
|
||||||
|
|
||||||
pgn = pl.when(pf < 240).then(id_shifted_8 % 65536).otherwise(id_shifted_8 % 262144).cast(pl.UInt32)
|
pgn = pl.when(pf < 240).then(
|
||||||
|
((id_int // 256) & 0x3FF00)
|
||||||
|
).otherwise(
|
||||||
|
((id_int // 256) & 0x3FFFF)
|
||||||
|
).cast(pl.UInt32)
|
||||||
|
|
||||||
return lf.with_columns(
|
return lf.with_columns(
|
||||||
pl.struct([
|
pl.struct([
|
||||||
@@ -44,26 +41,35 @@ def decode_j1939_metadata(lf: pl.LazyFrame) -> pl.LazyFrame:
|
|||||||
def decode_j1939_frames(df: pl.DataFrame) -> pl.DataFrame:
|
def decode_j1939_frames(df: pl.DataFrame) -> pl.DataFrame:
|
||||||
id_int = pl.col("ID").str.to_integer(base=16).cast(pl.UInt32)
|
id_int = pl.col("ID").str.to_integer(base=16).cast(pl.UInt32)
|
||||||
|
|
||||||
is_j1939 = (id_int > 0x7FF) & ((id_int % 33554432 // 16777216) == 0) & (pl.col("DLC") <= 8)
|
is_j1939 = id_int > 0x7FF
|
||||||
|
|
||||||
priority = (id_int // 67108864) % 8
|
priority = ((id_int // 67108864) % 8).cast(pl.UInt8)
|
||||||
pf = (id_int // 65536) % 256
|
pf = ((id_int // 65536) % 256).cast(pl.UInt8)
|
||||||
ps = (id_int // 256) % 256
|
ps = ((id_int // 256) % 256).cast(pl.UInt8)
|
||||||
sa = id_int % 256
|
sa = (id_int % 256).cast(pl.UInt8)
|
||||||
da = pl.when(pf < 240).then(ps).otherwise(255)
|
|
||||||
pgn = pl.when(pf < 240).then((id_int // 256) % 65536).otherwise((id_int // 256) % 262144)
|
da = pl.when(pf < 240).then(ps).otherwise(pl.lit(255, dtype=pl.UInt8)).cast(pl.UInt8)
|
||||||
|
|
||||||
|
pgn = pl.when(pf < 240).then(
|
||||||
|
((id_int // 256) & 0x3FF00)
|
||||||
|
).otherwise(
|
||||||
|
((id_int // 256) & 0x3FFFF)
|
||||||
|
).cast(pl.UInt32)
|
||||||
|
|
||||||
j1939_meta = pl.when(is_j1939).then(
|
j1939_meta = pl.when(is_j1939).then(
|
||||||
pl.struct([
|
pl.struct([
|
||||||
priority.cast(pl.UInt8).alias("Priority"),
|
priority.alias("Priority"),
|
||||||
pf.cast(pl.UInt8).alias("PF"),
|
pf.alias("PF"),
|
||||||
ps.cast(pl.UInt8).alias("PS"),
|
ps.alias("PS"),
|
||||||
sa.cast(pl.UInt8).alias("SA"),
|
sa.alias("SA"),
|
||||||
da.cast(pl.UInt8).alias("DA"),
|
da.alias("DA"),
|
||||||
pgn.cast(pl.UInt32).alias("PGN")
|
pgn.alias("PGN")
|
||||||
])
|
])
|
||||||
).otherwise(None)
|
).otherwise(None)
|
||||||
|
|
||||||
|
return df.with_columns(j1939_meta.alias("j1939_metadata"))
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
parser = argparse.ArgumentParser(description="J1939 decoder")
|
parser = argparse.ArgumentParser(description="J1939 decoder")
|
||||||
parser.add_argument("input_parquet", help="Path to the raw .parquet file")
|
parser.add_argument("input_parquet", help="Path to the raw .parquet file")
|
||||||
@@ -72,4 +78,5 @@ if __name__ == "__main__":
|
|||||||
|
|
||||||
df = pl.scan_parquet(args.input_parquet).collect()
|
df = pl.scan_parquet(args.input_parquet).collect()
|
||||||
decoded_df = decode_j1939_frames(df)
|
decoded_df = decode_j1939_frames(df)
|
||||||
|
|
||||||
decoded_df.write_parquet(args.output_parquet)
|
decoded_df.write_parquet(args.output_parquet)
|
||||||
|
|||||||
@@ -0,0 +1,76 @@
|
|||||||
|
import argparse
|
||||||
|
import json
|
||||||
|
from pathlib import Path
|
||||||
|
import fastparquet
|
||||||
|
import pandas as pd
|
||||||
|
import plotly.express as px
|
||||||
|
import plotly.graph_objects as go
|
||||||
|
|
||||||
|
|
||||||
|
def _extract_identifier(row: pd.Series) -> str:
|
||||||
|
"""Extracts PGN from metadata or falls back to CAN ID."""
|
||||||
|
meta = row.get('j1939_metadata')
|
||||||
|
if pd.isna(meta):
|
||||||
|
return f"{row['ID']} "
|
||||||
|
if isinstance(meta, str):
|
||||||
|
try:
|
||||||
|
meta = json.loads(meta)
|
||||||
|
except json.JSONDecodeError:
|
||||||
|
return f"{row['ID']}"
|
||||||
|
if isinstance(meta, dict) and 'PGN' in meta:
|
||||||
|
return f"PGN: {meta['PGN']} "
|
||||||
|
return f"{row['ID']}"
|
||||||
|
|
||||||
|
|
||||||
|
def load_data(file_path: Path) -> pd.DataFrame:
|
||||||
|
"""Loads Parquet file and adds an Identifier column."""
|
||||||
|
df = pd.read_parquet(file_path)
|
||||||
|
df['Identifier'] = df.apply(_extract_identifier, axis=1)
|
||||||
|
return df
|
||||||
|
|
||||||
|
|
||||||
|
def calculate_frequency(df: pd.DataFrame) -> pd.DataFrame:
|
||||||
|
"""Calculates frequency counts and percentages for identifiers."""
|
||||||
|
freq_df = df['Identifier'].value_counts().reset_index()
|
||||||
|
freq_df.columns = ['Identifier', 'Count']
|
||||||
|
total_messages = freq_df['Count'].sum()
|
||||||
|
freq_df['Percentage'] = (freq_df['Count'] / total_messages * 100).round(2)
|
||||||
|
return freq_df.sort_values('Count', ascending=True)
|
||||||
|
|
||||||
|
|
||||||
|
def visualize_frequency(stats_df: pd.DataFrame, title: str = "CAN Bus Message Frequency") -> go.Figure:
|
||||||
|
"""Generates an interactive Plotly horizontal bar chart with a logarithmic x-axis."""
|
||||||
|
fig = px.bar(
|
||||||
|
stats_df, y='Identifier', x='Count', orientation='h', title=title, log_x=True,
|
||||||
|
labels={'Identifier': 'PGN / CAN ID', 'Count': 'Message Count'},
|
||||||
|
color='Count', color_continuous_scale='Turbo',
|
||||||
|
hover_data={'Percentage': ':.2f', 'Count': True, 'Identifier': True}
|
||||||
|
)
|
||||||
|
|
||||||
|
fig.update_layout(
|
||||||
|
height=max(600, len(stats_df) * 18), width=1000,
|
||||||
|
xaxis_title='Total Message Count (Log Scale)', yaxis_title='PGN or CAN ID',
|
||||||
|
yaxis={'categoryorder': 'total ascending'},
|
||||||
|
plot_bgcolor='rgba(0,0,0,0)', paper_bgcolor='white',
|
||||||
|
font=dict(family="Segoe UI, Arial, sans-serif", size=12),
|
||||||
|
hoverlabel=dict(bgcolor="white", font_size=13, font_family="Segoe UI"),
|
||||||
|
margin=dict(l=250, r=50, t=80, b=50),
|
||||||
|
title=dict(font=dict(size=20), x=0.5)
|
||||||
|
)
|
||||||
|
|
||||||
|
fig.update_traces(
|
||||||
|
hovertemplate="<b>%{y}</b><br>Count: %{x:,}<br>Share: %{customdata[0]}%<extra></extra>"
|
||||||
|
)
|
||||||
|
return fig
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
parser = argparse.ArgumentParser(description="Analyze CAN bus Parquet data.")
|
||||||
|
parser.add_argument("-i", "--input", type=Path, required=True)
|
||||||
|
parser.add_argument("-o", "--output", type=Path, default=Path("can_analysis_report.html"))
|
||||||
|
parser.add_argument("-t", "--title", type=str, default="CAN Bus Message Frequency by PGN / ID")
|
||||||
|
args = parser.parse_args()
|
||||||
|
|
||||||
|
df = load_data(args.input)
|
||||||
|
stats_df = calculate_frequency(df)
|
||||||
|
fig = visualize_frequency(stats_df, title=args.title)
|
||||||
|
fig.write_html(str(args.output), include_plotlyjs='cdn')
|
||||||
Reference in New Issue
Block a user