0Pricing
AI Agents · Урок

Обработка временных рядов агентами

Скользящие окна, агрегация и обнаружение аномалий в потоковых данных датчиков.

«Обработка временных рядов агентами» — бесплатный урок 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 — локальная установка не требуется.

Все уроки этого курса

  1. Протокол MQTT для интеграции агентов
  2. Обработка временных рядов агентами
  3. Автоматическая реакция на события датчиков
  4. Развёртывание лёгких агентов на периферии
← Назад к AI Agents