Elaborazione dei dati di serie temporali negli agenti
Finestre mobili, aggregazione e rilevamento delle anomalie sui dati dei sensori in streaming.
Elaborazione dei dati di serie temporali negli agenti è una lezione AI Agents gratuita su CoddyKit. Questa è la lezione 2 di 4. Puoi leggere la lezione completa qui gratuitamente — poi esercitati direttamente nel browser con un editor di codice integrato e un tutor IA disponibile 24/7. Fa parte del percorso di apprendimento AI Agents, e i tuoi progressi si sincronizzano tra il web e l'app CoddyKit. Il corso AI Agents include 4 lezioni in totale.
Serie temporali negli agenti IoT
I dati dei sensori arrivano sotto forma di una serie temporale: una sequenza di coppie (timestamp, valore). Le letture grezze dei sensori contengono rumore, lacune e picchi occasionali. Gli agenti che agiscono sui dati grezzi senza elaborarli spesso generano falsi allarmi o non rilevano eventi reali.
L'elaborazione delle serie temporali trasforma i segnali grezzi in informazioni utili per agire.
Creazione di un buffer a finestra mobile
Una finestra mobile conserva in memoria solo le ultime N letture. Quando la finestra è piena, la lettura più vecchia viene eliminata all'arrivo di quella più recente. Questo è il fondamento di tutte le analisi delle serie temporali negli agenti.
from collections import deque
from datetime import datetime
class SensorBuffer:
def __init__(self, window_size: int = 60):
self.window_size = window_size
self._data = deque(maxlen=window_size)
def add(self, value: float, timestamp: datetime = None):
ts = timestamp or datetime.utcnow()
self._data.append({'ts': ts, 'value': value})
def values(self) -> list:
return [d['value'] for d in self._data]
def timestamps(self) -> list:
return [d['ts'] for d in self._data]
def is_full(self) -> bool:
return len(self._data) == self.window_size
buf = SensorBuffer(window_size=60)
buf.add(22.5)
buf.add(22.8)
buf.add(23.1)
print(f'Buffer: {len(buf._data)} readings, values: {buf.values()}')Media mobile
Una media mobile semplice (SMA) riduce il rumore calcolando la media degli ultimi N valori. Riduce l'effetto dei singoli malfunzionamenti dei sensori e rende visibile la tendenza sottostante. La utilizzi come riferimento per il rilevamento delle anomalie.
import statistics
def simple_moving_average(values: list, window: int = 10) -> list:
if len(values) < window:
return []
return [
statistics.mean(values[i - window:i])
for i in range(window, len(values) + 1)
]
def exponential_moving_average(values: list, alpha: float = 0.2) -> list:
"""EMA weights recent values more heavily."""
if not values:
return []
ema = [values[0]]
for v in values[1:]:
ema.append(alpha * v + (1 - alpha) * ema[-1])
return ema
readings = [22.1, 22.3, 22.0, 35.0, 22.2, 22.4, 22.1, 22.3, 22.5, 22.2, 22.4]
sma = simple_moving_average(readings, window=5)
ema = exponential_moving_average(readings, alpha=0.2)
print(f'SMA (last 3): {[round(v,2) for v in sma[-3:]]}')
print(f'EMA (last 3): {[round(v,2) for v in ema[-3:]]}')Rilevamento dei picchi
Un picco è una lettura che si discosta dalla tendenza recente di oltre N deviazioni standard (metodo dello z-score). Questo è l'approccio più comune al rilevamento delle anomalie nei dati dei sensori.
import statistics
def detect_spikes(
values: list,
window: int = 20,
z_threshold: float = 3.0
) -> list:
"""Returns list of (index, value, z_score) for detected spikes."""
if len(values) < window:
return []
spikes = []
for i in range(window, len(values)):
window_vals = values[i - window:i]
mean = statistics.mean(window_vals)
stdev = statistics.stdev(window_vals)
if stdev == 0:
continue
z_score = abs(values[i] - mean) / stdev
if z_score > z_threshold:
spikes.append({
'index': i,
'value': values[i],
'z_score': round(z_score, 2),
'mean': round(mean, 2)
})
return spikes
data = [22.1, 22.3, 22.0, 22.2, 22.4] * 5 + [55.0] + [22.2, 22.3] * 3
spikes = detect_spikes(data, window=10, z_threshold=3.0)
print('Spikes detected:', spikes)Pandas per l'analisi delle serie temporali
Per un'analisi più sofisticata, carichi i dati dei sensori in un DataFrame pandas con un DatetimeIndex. Pandas offre operazioni integrate di finestra mobile, ricampionamento e interpolazione, molto più rapide da scrivere rispetto ai cicli manuali.
import pandas as pd
from datetime import datetime, timedelta
# Create a sample time series DataFrame
base_time = datetime(2024, 1, 1, 12, 0, 0)
times = [base_time + timedelta(seconds=i*10) for i in range(20)]
values = [22.1, 22.3, None, 22.0, 22.4, 22.2, 35.0, 22.1,
22.3, 22.2, 22.5, 22.1, None, 22.4, 22.2, 22.3,
22.0, 22.1, 22.4, 22.2]
df = pd.DataFrame({'value': values}, index=pd.DatetimeIndex(times))
df.index.name = 'timestamp'
print('Shape:', df.shape)
print('Missing values:', df['value'].isna().sum())
print(df.head())Gestione dei timestamp mancanti
Le reti di sensori spesso perdono letture a causa di problemi di connettività. Pandas è in grado di rilevare e colmare le lacune: resample crea una griglia regolare, mentre interpolate riempie i valori mancanti linearmente o con il riempimento in avanti. Registri sempre nel log quanti valori sono stati imputati.
import pandas as pd
def fill_missing_readings(df: pd.DataFrame, freq: str = '10S') -> pd.DataFrame:
"""
df: DataFrame with DatetimeIndex and 'value' column
freq: expected sampling frequency ('10S' = 10 seconds, '1T' = 1 minute)
"""
original_count = df['value'].notna().sum()
# Resample to regular grid (introduces NaN for missing periods)
df_regular = df.resample(freq).mean()
missing_count = df_regular['value'].isna().sum()
print(f'Missing readings before fill: {missing_count}')
# Forward fill then linear interpolate
df_regular['value'] = df_regular['value'].interpolate(
method='linear', limit=5 # don't fill gaps longer than 5 periods
)
filled_count = df_regular['value'].notna().sum()
print(f'Filled {filled_count - original_count} missing values')
return df_regularRicampionamento: aggregazione da 1 minuto a 5 minuti
Il ricampionamento riduce la frequenza dei dati ad alta frequenza, portandoli a una risoluzione più grossolana. In questo modo diminuiscono il rumore e i requisiti di archiviazione. Utilizzi resample('5T').agg() per calcolare minimo, massimo, media e deviazione standard su ogni finestra di 5 minuti.
import pandas as pd
def resample_to_5min(df: pd.DataFrame) -> pd.DataFrame:
return df.resample('5min').agg({
'value': ['mean', 'min', 'max', 'std', 'count']
}).round(3)
# Example with 1-minute data:
times = pd.date_range('2024-01-01 12:00', periods=30, freq='1min')
import random
random.seed(42)
vals = [22.0 + random.gauss(0, 0.5) for _ in range(30)]
df_1min = pd.DataFrame({'value': vals}, index=times)
df_5min = resample_to_5min(df_1min)
print(df_5min)Rilevamento delle tendenze con la regressione lineare
La temperatura sta aumentando costantemente oppure presenta solo rumore? Esegua una regressione lineare sulla finestra mobile. Una pendenza positiva indica una tendenza al rialzo; una pendenza che supera una soglia attiva un avviso prima ancora che la soglia venga raggiunta.
def detect_trend(
values: list,
slope_threshold: float = 0.1 # units per second
) -> dict:
import statistics
n = len(values)
if n < 2:
return {'trend': 'insufficient_data'}
x = list(range(n))
x_mean = statistics.mean(x)
y_mean = statistics.mean(values)
numerator = sum((xi - x_mean) * (yi - y_mean) for xi, yi in zip(x, values))
denominator = sum((xi - x_mean) ** 2 for xi in x)
slope = numerator / denominator if denominator != 0 else 0.0
return {
'slope': round(slope, 4),
'trend': 'rising' if slope > slope_threshold
else 'falling' if slope < -slope_threshold
else 'stable',
'alert': abs(slope) > slope_threshold * 2
}
readings = [22.0, 22.5, 23.0, 23.5, 24.0, 24.5, 25.0]
print(detect_trend(readings, slope_threshold=0.3))Decisioni dell'agente basate sulle serie temporali
Combini il rilevamento dei picchi, il rilevamento delle tendenze e la media mobile per prendere una decisione composita. L'agente utilizza una gerarchia di regole: i picchi attivano azioni immediate, le tendenze attivano avvisi e l'LLM gestisce i casi ambigui.
def analyze_sensor_window(
values: list,
topic: str
) -> dict:
if len(values) < 10:
return {'action': 'collecting_data'}
spikes = detect_spikes(values, window=10, z_threshold=3.0)
trend_info = detect_trend(values[-20:], slope_threshold=0.2)
avg = sum(values[-10:]) / 10
if spikes:
return {
'action': 'IMMEDIATE_ALERT',
'reason': f'Spike detected: {spikes[-1]["value"]} (z={spikes[-1]["z_score"]})',
'severity': 'high'
}
if trend_info['alert']:
return {
'action': 'TREND_WARNING',
'reason': f'Rapid {trend_info["trend"]} trend: slope={trend_info["slope"]}',
'severity': 'medium'
}
return {
'action': 'NORMAL',
'avg_last_10': round(avg, 2),
'trend': trend_info['trend']
}
result = analyze_sensor_window([22.0]*15 + [55.0], 'sensors/temp')
print(result)Persistenza delle serie temporali in un database
Per le analisi a lungo termine, memorizzi le letture dei sensori in un database per serie temporali. TimescaleDB (un'estensione di PostgreSQL) e InfluxDB sono scelte diffuse. L'utilizzo di psycopg2 con TimescaleDB consente di interrogare i dati con SQL standard e funzioni specifiche per il tempo.
import psycopg2
from datetime import datetime
# TimescaleDB connection (standard PostgreSQL connection)
conn = psycopg2.connect(
host='localhost', port=5432, dbname='iot',
user='agent', password='YOUR_DB_PASSWORD'
)
def insert_reading(topic: str, value: float, ts: datetime = None):
ts = ts or datetime.utcnow()
with conn.cursor() as cur:
cur.execute(
'INSERT INTO sensor_readings (time, topic, value) VALUES (%s, %s, %s)',
(ts, topic, value)
)
conn.commit()
def query_last_hour(topic: str) -> list:
with conn.cursor() as cur:
cur.execute(
'SELECT time, value FROM sensor_readings '
'WHERE topic=%s AND time > NOW() - INTERVAL \'1 hour\' '
'ORDER BY time ASC',
(topic,)
)
return cur.fetchall()Periodo di cooldown degli avvisi
Senza un cooldown, un'anomalia persistente genera centinaia di avvisi al minuto. Implementi un cooldown per topic: dopo l'invio di un avviso per un topic, sopprima gli ulteriori avvisi per quel topic per N secondi.
from datetime import datetime, timedelta
class AlertCooldownManager:
def __init__(self, cooldown_seconds: int = 300):
self.cooldown_seconds = cooldown_seconds
self._last_alert: dict = {} # topic -> last alert datetime
def should_alert(self, topic: str) -> bool:
last = self._last_alert.get(topic)
if last is None:
return True
return (datetime.utcnow() - last).seconds >= self.cooldown_seconds
def mark_alerted(self, topic: str):
self._last_alert[topic] = datetime.utcnow()
def cooldown_remaining(self, topic: str) -> int:
last = self._last_alert.get(topic)
if last is None:
return 0
elapsed = (datetime.utcnow() - last).seconds
return max(0, self.cooldown_seconds - elapsed)
cooldown = AlertCooldownManager(cooldown_seconds=300)
if cooldown.should_alert('sensors/temperature'):
print('Sending alert')
cooldown.mark_alerted('sensors/temperature')
else:
print(f'Cooldown: {cooldown.cooldown_remaining("sensors/temperature")}s remaining')Verifica delle conoscenze
Che cosa indica uno z-score superiore a 3,0 quando viene utilizzato per rilevare i picchi?
Riepilogo: elaborazione dei dati delle serie temporali
Ha esaminato l'intero kit di strumenti per l'elaborazione delle serie temporali negli agenti:
- Buffer a finestra mobile: deque con maxlen per uno streaming efficiente in termini di memoria
- Media mobile: SMA ed EMA per la riduzione del rumore
- Rilevamento dei picchi: metodo dello z-score applicato alle statistiche della finestra mobile
- Pandas: ricampionamento, interpolazione e aggregazione con DatetimeIndex
- Rilevamento delle tendenze: pendenza della regressione lineare per un allarme preventivo
- Cooldown degli avvisi: soppressione degli avvisi ripetuti per anomalie persistenti
Prossimo argomento: risposte automatizzate degli agenti agli eventi dei sensori — code di azioni, deduplicazione e criteri per le azioni.
Domande Frequenti
La lezione «Elaborazione dei dati di serie temporali negli agenti» è gratuita?
Sì — il testo completo di «Elaborazione dei dati di serie temporali negli agenti» è gratuito qui sul web. Per esercitarvi in modo interattivo (un editor di codice integrato e un tutor IA 24/7) e sbloccare il resto del corso AI Agents, passa a CoddyKit PRO. Il corso AI Agents include 4 lezioni in totale.
Cosa imparerò in «Elaborazione dei dati di serie temporali negli agenti»?
Finestre mobili, aggregazione e rilevamento delle anomalie sui dati dei sensori in streaming. Eserciti AI Agents con codice pratico che esegui direttamente nel browser, e un tutor IA 24/7 risponde alle tue domande mentre lavori sulla lezione.
Ho bisogno di esperienza per iniziare AI Agents?
Non è richiesta alcuna esperienza precedente. AI Agents su CoddyKit è strutturato per principianti e studenti avanzati, quindi puoi iniziare da qui o dall'inizio e procedere al tuo ritmo. Questa è la lezione 2 di 4.
Quanto tempo richiede la lezione «Elaborazione dei dati di serie temporali negli agenti»?
La maggior parte delle lezioni CoddyKit richiede circa 5–10 minuti. Ogni lezione è breve e interattiva, quindi fai progressi costanti e riprendi esattamente da dove hai lasciato su web e app.
Posso scrivere ed eseguire codice in questa lezione AI Agents?
Sì. Ogni lezione AI Agents include un editor di codice integrato, quindi scrivi ed esegui codice reale direttamente nel tuo browser e ricevi feedback istantaneo dall'IA — nessuna configurazione locale necessaria.
Tutte le lezioni di questo corso
- Protocollo MQTT per l’integrazione degli agenti
- Elaborazione dei dati di serie temporali negli agenti
- Risposta automatizzata agli eventi dei sensori
- Distribuzione edge di agenti leggeri