Обработка временных рядов агентами
Скользящие окна, агрегация и обнаружение аномалий в потоковых данных датчиков.
«Обработка временных рядов агентами» — бесплатный урок AI Agents на CoddyKit. Это урок 2 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения AI Agents, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс AI Agents содержит 4 уроков всего.
Временные ряды в агентах IoT
Данные датчиков поступают в виде временного ряда: последовательности пар (временная метка, значение). Исходные показания датчиков содержат шум, пропуски и редкие выбросы. Агенты, которые действуют на основе необработанных данных без предварительной обработки, часто вызывают ложные тревоги или пропускают реальные события.
Обработка временных рядов преобразует исходные сигналы в практические выводы.
Создание буфера скользящего окна
Скользящее окно хранит в памяти только последние N показаний. Когда окно заполнено, при поступлении нового показания самое старое удаляется. Это основа всего анализа временных рядов в агентах.
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()}')Скользящее среднее
Простое скользящее среднее (SMA) сглаживает шум, усредняя последние N значений. Оно уменьшает влияние отдельных сбоев датчика и показывает основной тренд. Используйте его как базовый уровень для обнаружения аномалий.
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:]]}')Обнаружение выбросов
Выброс — это показание, которое отклоняется от недавнего тренда более чем на N стандартных отклонений (метод z-оценки). Это самый распространённый подход к обнаружению аномалий в данных датчиков.
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 для анализа временных рядов
Для более сложного анализа загрузите данные датчиков в pandas DataFrame с DatetimeIndex. Pandas предоставляет встроенные операции скользящего окна, пересэмплирования и интерполяции, которые позволяют написать решение значительно быстрее, чем при использовании ручных циклов.
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())Обработка пропущенных временных меток
В сетях датчиков показания часто пропускаются из-за проблем с подключением. Pandas умеет обнаруживать и заполнять пропуски: resample создаёт регулярную сетку, а interpolate заполняет пропущенные значения линейно или методом прямого заполнения. Всегда записывайте в журнал количество восстановленных значений.
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_regularПересэмплирование: агрегация с одной до пяти минут
Пересэмплирование уменьшает частоту данных, переводя их к более грубому разрешению. Это снижает уровень шума и требования к объёму хранилища. Используйте resample('5T').agg(), чтобы вычислить минимум, максимум, среднее и стандартное отклонение для каждого пятиминутного окна.
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)Обнаружение тренда с помощью линейной регрессии
Температура постоянно повышается или это просто шум? Выполните линейную регрессию по скользящему окну. Положительный наклон указывает на восходящий тренд, а наклон, превышающий порог, вызывает оповещение ещё до достижения самого порога.
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))Принятие решений агентом по временному ряду
Объедините обнаружение выбросов, обнаружение тренда и скользящее среднее, чтобы принять комплексное решение. Агент использует иерархию правил: выбросы запускают немедленные действия, тренды вызывают предупреждения, а неоднозначные случаи обрабатывает LLM.
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)Сохранение временных рядов в базе данных
Для долгосрочного анализа сохраняйте показания датчиков в базе данных временных рядов. TimescaleDB (расширение PostgreSQL) и InfluxDB — популярные варианты. Использование psycopg2 с TimescaleDB позволяет запрашивать данные с помощью стандартного SQL и специализированных функций для работы со временем.
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()Период охлаждения для оповещений
Без периода охлаждения длительная аномалия вызывает сотни оповещений в минуту. Реализуйте отдельный период охлаждения для каждого топика: после отправки оповещения для топика подавляйте дальнейшие оповещения по нему в течение N секунд.
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')Проверка знаний
Что означает z-оценка выше 3,0 при обнаружении выбросов?
Итоги: обработка данных временных рядов
Вы изучили полный набор средств обработки временных рядов для агентов:
- Буфер скользящего окна: deque с maxlen для эффективной работы с памятью при потоковой обработке
- Скользящее среднее: SMA и EMA для снижения шума
- Обнаружение выбросов: метод z-оценки на основе статистики скользящего окна
- Pandas: пересэмплирование, интерполяция и агрегация с DatetimeIndex
- Обнаружение тренда: наклон линейной регрессии для раннего предупреждения
- Период охлаждения для оповещений: подавление повторных оповещений при длительных аномалиях
Далее: автоматические ответы агента на события датчиков — очереди действий, дедупликация и политики действий.
Часто задаваемые вопросы
Урок «Обработка временных рядов агентами» бесплатный?
Да — полный текст урока «Обработка временных рядов агентами» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс AI Agents, подпишись на CoddyKit PRO. Курс AI Agents содержит 4 уроков всего.
Чему я научусь в уроке «Обработка временных рядов агентами»?
Скользящие окна, агрегация и обнаружение аномалий в потоковых данных датчиков. Ты практикуешь AI Agents с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.
Нужен ли мне опыт, чтобы начать AI Agents?
Предыдущий опыт не требуется. AI Agents на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 2 из 4.
Сколько времени занимает урок «Обработка временных рядов агентами»?
Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.
Можно ли писать и запускать код в этом уроке AI Agents?
Да. Каждый урок AI Agents включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.
Все уроки этого курса
- Протокол MQTT для интеграции агентов
- Обработка временных рядов агентами
- Автоматическая реакция на события датчиков
- Развёртывание лёгких агентов на периферии