О проекте - Glost/iot-olap-clickhouse GitHub Wiki

Цель проекта - продемонстрировать возможности использования столбцовой OLAP-СУБД ClickHouse на примере IoT-системы с "датчиками", генерирующими данные.

Авторы: Александр Плесовских, Антон Ригин, Максим Шмаков.

Общая информация

Разработана система, эмулирующая множество датчиков, которые передают серверу (бэкенд-приложению) показатели. Данные показатели сохраняются в ClickHouse, чтобы затем получать по ним агрегированные данные и строить различные графики (последнее поддержано в мобильном приложении для iOS). В представленном прототипе системы датчики представлены генератором случайных значений (при этом генерируемые им показатели находятся в "разумных" границах для той или иной физической величины и соблюдается "плавность" изменения показаний).

Архитектура системы

Система состоит из пяти компонентов:

  • Столбцовая OLAP-СУБД ClickHouse - для долговременного хранения показаний датчиков и получения по ним агрегированных данных
  • Реляционная OLTP-СУБД PostgreSQL - для кратковременного (не более 1 часа) хранения показаний датчиков до их выгрузки в ClickHouse
  • Сервер (бэкенд-приложение), обеспечивающее приём показаний от датчиков, их выгрузку в ClickHouse и "прослойку" между клиентами и ClickHouse для получения из ClickHouse агрегированных данных; технологический стек: Java 8, Spring Boot 2, ClickHouse JDBC (официальный JDBC-клиент ClickHouse), Liquibase (для PostgreSQL-миграций), nginx
  • Генератор показаний датчиков, технологический стек: NodeJS
  • Клиент (мобильное iOS-приложение), обеспечивающее запрос агрегированных данных через сервер, их получение и отображение построенных по ним графиков

Принципы взаимодействия компонентов системы следующие:

  • С обеими используемыми в данной работе СУБД (ClickHouse и PostgreSQL) напрямую взаимодействует только сервер (бэкенд-приложение), а генератор показаний датчиков и клиент (мобильное iOS-приложение) напрямую с этими СУБД не взаимодействуют (только через сервер)
  • Генератор показаний датчиков и клиент (мобильное iOS-приложение) напрямую взаимодействуют только с сервером (через REST API), а друг с другом они не взаимодействуют
  • Сервер (бэкенд-приложение) при этом взаимодействует напрямую со всеми остальными компонентами системы

Датчики и их показания

Генератором показаний датчиков эмулированы 15 датчиков, каждый со своими координатами.

Координаты датчиков фиксированы и задаются следующим образом:

  • latitude - широта (в градусах), на которой расположен датчик - тип double
  • longitude - долгота (в градусах), на которой расположен датчик - тип double
  • altitude - высота над уровнем моря (в метрах), на которой расположен датчик - тип double

Показания датчиков состоят из значений следующих величин:

  • air_temperature - температура воздуха (в градусах Цельсия) - тип double
  • air_humidity - влажность воздуха (в процентах) - тип double
  • wind_speed - скорость ветра (в метрах в секунду) - тип double
  • is_raining - флаг, идёт ли дождь (true, если идёт, иначе false) - тип boolean
  • illuminance - освещённостьлюксах) - тип double

Каждый из 15 датчиков передаёт показания раз в полминуты (30 секунд), таким образом, чтобы каждые 2 секунды свои показания передавал один и только один датчик (то есть у каждого из 15 датчиков есть фиксированная секунда, в рамках которой каждые полминуты он передаёт показания).

Вместе с показаниями датчик передаёт timestamp - время произведения замера показаний на датчике, с указанием часового пояса - строка (тип String) в формате yyyy-MM-dd['T'][ ]HH:mm:ss.SSSXXXX (например: 2019-12-01 17:18:48.126+0300, 2019-12-01T17:18:48.126+0300, 2019-12-01 17:18:48.126Z, 2019-12-01T17:18:48.126Z, и т. д.). Это необходимо для того, чтобы сервер имел именно то время показаний, в которое они были измерены датчиком, без учёта возможных задержек при передаче по сети.

Принципы работы сервера (бэкенд-приложения) и его взаимодействия с СУБД ClickHouse и PostgreSQL

Бэкенд-приложение работает в соответствии с описанным API.

Получение и временное хранение данных от датчиков

При каждом получении данных от датчика (с использованием метода POST /pushData) сервер их сохраняет для кратковременного (не более 1 часа) хранения в таблицу sensor_values в PostgreSQL, созданную с использованием следующего DDL-кода.

CREATE TABLE IF NOT EXISTS sensor_values (
    latitude DOUBLE PRECISION NOT NULL,
    longitude DOUBLE PRECISION NOT NULL,
    altitude DOUBLE PRECISION NOT NULL,
    timestamp TIMESTAMP NOT NULL,
    air_temperature DOUBLE PRECISION NOT NULL,
    air_humidity DOUBLE PRECISION NOT NULL,
    wind_speed DOUBLE PRECISION NOT NULL,
    is_raining BOOLEAN NOT NULL,
    illuminance DOUBLE PRECISION NOT NULL,
    PRIMARY KEY (latitude, longitude, altitude, timestamp)
);

Выгрузка данных в ClickHouse

Раз в час (в момент начала каждого астрономического часа) все данные из PostgreSQL-таблицы sensor_values выгружаются в ClickHouse и удаляются из PostgreSQL. Запись об этом добавляется в PostgreSQL-таблицу scheduled_pushing_data_to_clickhouse_log, созданную с использованием следующего DDL-кода.

CREATE TABLE IF NOT EXISTS scheduled_pushing_data_to_clickhouse_log (
    timestamp TIMESTAMP PRIMARY KEY,
    pushed_rows_count INTEGER NOT NULL,
    clickhouse_query_duration INTEGER NOT NULL
);

В ClickHouse данные сохраняются в таблицу sensor_values, созданную с использованием следующего DDL-кода.

CREATE TABLE IF NOT EXISTS sensor_values (
    latitude Float64,
    longitude Float64,
    altitude Float64,
    timestamp DateTime,
    air_temperature Float64,
    air_humidity Float64,
    wind_speed Float64,
    is_raining UInt8,
    illuminance Float64
) ENGINE = MergeTree()
PARTITION BY (toStartOfDay(timestamp), toStartOfHour(timestamp), toStartOfMinute(timestamp))
ORDER BY (timestamp, latitude, longitude, altitude);

MergeTree() в данном DDL-коде - это название соответствующего ClickHouse-движка таблицы, одного из самых популярных и функциональных среди движков таблиц ClickHouse, позволяющего, среди прочего:

ClickHouse не поддерживает тип Boolean, поэтому столбец is_raining имеет тип UInt8.

Выгрузка данных производится посредством следующего DML-выражения INSERT INTO.

INSERT INTO sensor_values (latitude,
                           longitude,
                           altitude,
                           timestamp,
                           air_temperature,
                           air_humidity,
                           wind_speed,
                           is_raining,
                           illuminance)

Вместе с этим выражением передаются сами данные. ClickHouse имеет особенность, согласно которой данные, передаваемые после INSERT INTO парсятся специальным быстрым потоковым парсером, повышая скорость исполнения запроса на вставку данных. Впрочем, в случае нашего приложения данные при вставке и так передаются в формате RowBinary, то есть уже в распарсенном виде.

Получение агрегированных данных из ClickHouse

Получение агрегированных данных из ClickHouse на сервере производится по запросу от клиента.

Получение агрегированных данных из ClickHouse, разделённых по отдельным датчикам

Для получения агрегированных данных из ClickHouse, разделённых по отдельным датчикам, используется SQL-запрос к ClickHouse, формируемый сервером из следующего шаблона.

SELECT latitude, longitude, altitude, timestamp_rounded,
       min_air_temperature,
       max_air_temperature,
       avg_air_temperature,
       median_air_temperature,
       var_air_temperature,
       min_air_humidity,
       max_air_humidity,
       avg_air_humidity,
       median_air_humidity,
       var_air_humidity,
       min_wind_speed,
       max_wind_speed,
       avg_wind_speed,
       median_wind_speed,
       var_wind_speed,
       min_illuminance,
       max_illuminance,
       avg_illuminance,
       median_illuminance,
       var_illuminance,
       avg_rains
FROM (SELECT latitude, longitude, altitude, timestamp_rounded,
             min(air_temperature) AS min_air_temperature,
             max(air_temperature) AS max_air_temperature,
             avg(air_temperature) AS avg_air_temperature,
             median(air_temperature) AS median_air_temperature,
             varPop(air_temperature) AS var_air_temperature,
             min(air_humidity) AS min_air_humidity,
             max(air_humidity) AS max_air_humidity,
             avg(air_humidity) AS avg_air_humidity,
             median(air_humidity) AS median_air_humidity,
             varPop(air_humidity) AS var_air_humidity,
             min(wind_speed) AS min_wind_speed,
             max(wind_speed) AS max_wind_speed,
             avg(wind_speed) AS avg_wind_speed,
             median(wind_speed) AS median_wind_speed,
             varPop(wind_speed) AS var_wind_speed,
             min(illuminance) AS min_illuminance,
             max(illuminance) AS max_illuminance,
             avg(illuminance) AS avg_illuminance,
             median(illuminance) AS median_illuminance,
             varPop(illuminance) AS var_illuminance,
             toFloat64(sum(is_raining)) AS avg_rains
      FROM sensor_values
      WHERE timestamp BETWEEN ? AND ? ${coordinate_intervals}
      GROUP BY latitude, longitude, altitude, ${timestamp_rounding_function}(timestamp) AS timestamp_rounded
      UNION ALL
      SELECT NULL AS latitude, NULL AS longitude, NULL AS altitude, timestamp_rounded,
          min(air_temperature) AS min_air_temperature,
          max(air_temperature) AS max_air_temperature,
          avg(air_temperature) AS avg_air_temperature,
          median(air_temperature) AS median_air_temperature,
          varPop(air_temperature) AS var_air_temperature,
          min(air_humidity) AS min_air_humidity,
          max(air_humidity) AS max_air_humidity,
          avg(air_humidity) AS avg_air_humidity,
          median(air_humidity) AS median_air_humidity,
          varPop(air_humidity) AS var_air_humidity,
          min(wind_speed) AS min_wind_speed,
          max(wind_speed) AS max_wind_speed,
          avg(wind_speed) AS avg_wind_speed,
          median(wind_speed) AS median_wind_speed,
          varPop(wind_speed) AS var_wind_speed,
          min(illuminance) AS min_illuminance,
          max(illuminance) AS max_illuminance,
          avg(illuminance) AS avg_illuminance,
          median(illuminance) AS median_illuminance,
          varPop(illuminance) AS var_illuminance,
          (SELECT avg(rains) FROM (SELECT sum(is_raining) AS rains
                                   FROM sensor_values
                                   WHERE timestamp BETWEEN ? AND ? ${coordinate_intervals}
                                   GROUP BY latitude, longitude, altitude,
                                   ${timestamp_rounding_function}(timestamp))) AS avg_rains
      FROM (SELECT timestamp, air_temperature, air_humidity, wind_speed, is_raining, illuminance
            FROM sensor_values
            WHERE timestamp BETWEEN ? AND ? ${coordinate_intervals})
      GROUP BY ${timestamp_rounding_function}(timestamp) AS timestamp_rounded)
ORDER BY latitude, longitude, altitude, timestamp_rounded;

Вторая часть запроса (после выражения UNION ALL) здесь позволяет дополнительно получить агрегированные по всем подпадающим по условия датчикам значения, без деления информации по отдельным датчикам.

Здесь плейсхолдер ${coordinate_intervals} заменяется интервалами координат датчиков, если таковые были переданы, а плейсхолдер ${timestamp_rounding_function} заменяется одной из функций - toStartOfDay, toStartOfHour либо toStartOfMinute, в зависимости от того, что указано в поле aggregated_period (DAY, HOUR либо MINUTE) запроса.

После этого формируется prepared statement, к которому добавляются значения параметров, обозначенных в запросе символами ?, после чего он обрабатывается ClickHouse JDBC, передаётся в ClickHouse, исполняется ClickHouse, на сервер возвращается ответ, который сервер преобразует в формат JSON, определённый в API, и возвращает клиенту.

Например, от клиента был получен запрос, показанный в примере в API (но чуть с другими значениями timestamp_interval - имеющимися в данных, сгенерированных в рамках нашего проекта), а именно следующий.

{
  "latitude_intervals": [
    {
      "from": 55.755831,
      "to": 57.854573
    },
    {
      "from": 65.459327,
      "to": 69.423342
    }
  ],
  "longitude_intervals": [
    {
      "from": 37.617673,
      "to": 41.965782
    }
  ],
  "altitude_intervals": [
    {
      "from": 150.0,
      "to": 450.0
    }
  ],
  "timestamp_interval": {
    "from": "2019-12-18 00:00:00.000+0300",
    "to": "2019-12-19 00:00:00.000+0300"
  },
  "aggregated_period": "HOUR",
  "need_split_by_coordinates": true
}

"need_split_by_coordinates": true в этом запросе как раз означает, что агрегированные данные необходимо делить по отдельным датчикам.

В таком случае SQL-запрос к ClickHouse будет выглядеть следующим образом (если допустить, что сервер и кластер ClickHouse работают в одном часовом поясе, в противном случае время будет сконвертировано в часовой пояс кластера ClickHouse, и если допустить, что мы заменили символы ? на значения без использования prepared statement, хотя это на деле это будет не так).

SELECT latitude, longitude, altitude, timestamp_rounded,
       min_air_temperature,
       max_air_temperature,
       avg_air_temperature,
       median_air_temperature,
       var_air_temperature,
       min_air_humidity,
       max_air_humidity,
       avg_air_humidity,
       median_air_humidity,
       var_air_humidity,
       min_wind_speed,
       max_wind_speed,
       avg_wind_speed,
       median_wind_speed,
       var_wind_speed,
       min_illuminance,
       max_illuminance,
       avg_illuminance,
       median_illuminance,
       var_illuminance,
       avg_rains
FROM (SELECT latitude, longitude, altitude, timestamp_rounded,
             min(air_temperature) AS min_air_temperature,
             max(air_temperature) AS max_air_temperature,
             avg(air_temperature) AS avg_air_temperature,
             median(air_temperature) AS median_air_temperature,
             varPop(air_temperature) AS var_air_temperature,
             min(air_humidity) AS min_air_humidity,
             max(air_humidity) AS max_air_humidity,
             avg(air_humidity) AS avg_air_humidity,
             median(air_humidity) AS median_air_humidity,
             varPop(air_humidity) AS var_air_humidity,
             min(wind_speed) AS min_wind_speed,
             max(wind_speed) AS max_wind_speed,
             avg(wind_speed) AS avg_wind_speed,
             median(wind_speed) AS median_wind_speed,
             varPop(wind_speed) AS var_wind_speed,
             min(illuminance) AS min_illuminance,
             max(illuminance) AS max_illuminance,
             avg(illuminance) AS avg_illuminance,
             median(illuminance) AS median_illuminance,
             varPop(illuminance) AS var_illuminance,
             toFloat64(sum(is_raining)) AS avg_rains
      FROM sensor_values
      WHERE timestamp BETWEEN '2019-12-18 00:00:00' AND '2019-12-19 00:00:00'
          AND (latitude BETWEEN 55.755831 AND 57.854573
              OR latitude BETWEEN 65.459327 AND 69.423342)
          AND (longitude BETWEEN 37.617673 AND 41.965782)
          AND (altitude BETWEEN 150.0 AND 450.0)
      GROUP BY latitude, longitude, altitude, toStartOfHour(timestamp) AS timestamp_rounded
      UNION ALL
      SELECT NULL AS latitude, NULL AS longitude, NULL AS altitude, timestamp_rounded,
          min(air_temperature) AS min_air_temperature,
          max(air_temperature) AS max_air_temperature,
          avg(air_temperature) AS avg_air_temperature,
          median(air_temperature) AS median_air_temperature,
          varPop(air_temperature) AS var_air_temperature,
          min(air_humidity) AS min_air_humidity,
          max(air_humidity) AS max_air_humidity,
          avg(air_humidity) AS avg_air_humidity,
          median(air_humidity) AS median_air_humidity,
          varPop(air_humidity) AS var_air_humidity,
          min(wind_speed) AS min_wind_speed,
          max(wind_speed) AS max_wind_speed,
          avg(wind_speed) AS avg_wind_speed,
          median(wind_speed) AS median_wind_speed,
          varPop(wind_speed) AS var_wind_speed,
          min(illuminance) AS min_illuminance,
          max(illuminance) AS max_illuminance,
          avg(illuminance) AS avg_illuminance,
          median(illuminance) AS median_illuminance,
          varPop(illuminance) AS var_illuminance,
          (SELECT avg(rains) FROM (SELECT sum(is_raining) AS rains
          FROM sensor_values
          WHERE timestamp BETWEEN '2019-12-18 00:00:00' AND '2019-12-19 00:00:00'
              AND (latitude BETWEEN 55.755831 AND 57.854573
                  OR latitude BETWEEN 65.459327 AND 69.423342)
              AND (longitude BETWEEN 37.617673 AND 41.965782)
              AND (altitude BETWEEN 150.0 AND 450.0)
          GROUP BY latitude, longitude, altitude,
                   toStartOfHour(timestamp))) AS avg_rains
      FROM (SELECT timestamp, air_temperature, air_humidity, wind_speed, is_raining, illuminance
          FROM sensor_values
            WHERE timestamp BETWEEN '2019-12-18 00:00:00' AND '2019-12-19 00:00:00'
                AND (latitude BETWEEN 55.755831 AND 57.854573
                    OR latitude BETWEEN 65.459327 AND 69.423342)
                AND (longitude BETWEEN 37.617673 AND 41.965782)
                AND (altitude BETWEEN 150.0 AND 450.0))
      GROUP BY toStartOfHour(timestamp) AS timestamp_rounded)
ORDER BY latitude, longitude, altitude, timestamp_rounded;

Необходимо отметить, что ClickHouse принимает и хранит данные о времени с точностью до секунд, данные о миллисекундах отбрасываются.

По умолчанию в официальном ClickHouse JDBC результат запроса запрашивается в формате TabSeparatedWithNamesAndTypes, в котором строки разделены символами переноса строк, а столбцы - символами табуляции. Результат парсится ClickHouse JDBC, благодаря чему получается объект типа ClickHouseResponse, в котором каждая ячейка записана как строка, и она уже парсится на сервере в вещественные числа, дату-время, и т. д., после чего формируется JSON-ответ для клиента, который ему и возвращается.

После выполнения приведённого выше запроса клиенту вернётся ответ в формате JSON, приведённый в файле split-response.json. Подробнее о значении разных полей ответа можно почитать в API: в разделах про пользовательские типы данных и про метод метод POST /getAggregatedData. Поле click_house_query_execution_time представляет количество миллисекунд, в течение которых выполнялся запрос на получение агрегированных данных из ClickHouse.

Получение агрегированных данных из ClickHouse, без деления информации по отдельным датчикам

Для получения агрегированных данных из ClickHouse, без деления информации по отдельным датчикам, используется SQL-запрос к ClickHouse, формируемый сервером из следующего шаблона.

SELECT timestamp_rounded,
       min(air_temperature) AS min_air_temperature,
       max(air_temperature) AS max_air_temperature,
       avg(air_temperature) AS avg_air_temperature,
       median(air_temperature) AS median_air_temperature,
       varPop(air_temperature) AS var_air_temperature,
       min(air_humidity) AS min_air_humidity,
       max(air_humidity) AS max_air_humidity,
       avg(air_humidity) AS avg_air_humidity,
       median(air_humidity) AS median_air_humidity,
       varPop(air_humidity) AS var_air_humidity,
       min(wind_speed) AS min_wind_speed,
       max(wind_speed) AS max_wind_speed,
       avg(wind_speed) AS avg_wind_speed,
       median(wind_speed) AS median_wind_speed,
       varPop(wind_speed) AS var_wind_speed,
       min(illuminance) AS min_illuminance,
       max(illuminance) AS max_illuminance,
       avg(illuminance) AS avg_illuminance,
       median(illuminance) AS median_illuminance,
       varPop(illuminance) AS var_illuminance,
       (SELECT avg(rains) FROM (SELECT sum(is_raining) AS rains
                                FROM sensor_values
                                WHERE timestamp BETWEEN ? AND ? ${coordinate_intervals}
                                GROUP BY latitude, longitude, altitude,
                                ${timestamp_rounding_function}(timestamp))) AS avg_rains
FROM (SELECT timestamp, air_temperature, air_humidity, wind_speed, is_raining, illuminance
      FROM sensor_values
      WHERE timestamp BETWEEN ? AND ? ${coordinate_intervals})
GROUP BY ${timestamp_rounding_function}(timestamp) AS timestamp_rounded
ORDER BY timestamp_rounded;

Здесь плейсхолдеры ${coordinate_intervals} и ${timestamp_rounding_function} имеют те же значения, которые описаны в предыдущем подразделе, а после формирования окончательного запроса используется prepared statement так же, как описано в предыдущем подразделе.

Например, от клиента был получен запрос, показанный в примере в API (но чуть с другими значениями timestamp_interval - имеющимися в данных, сгенерированных в рамках нашего проекта), а именно следующий.

{
  "altitude_intervals": [
    {
      "from": 150.0,
      "to": 450.0
    }
  ],
  "timestamp_interval": {
    "from": "2019-12-18 00:00:00.000+0300",
    "to": "2019-12-19 00:00:00.000+0300"
  },
  "aggregated_period": "HOUR",
  "need_split_by_coordinates": false
}

"need_split_by_coordinates": false в этом запросе как раз означает, что агрегированные данные не нужно делить по отдельным датчикам.

В таком случае SQL-запрос к ClickHouse будет выглядеть следующим образом (если допустить, что сервер и кластер ClickHouse работают в одном часовом поясе, в противном случае время будет сконвертировано в часовой пояс кластера ClickHouse, и если допустить, что мы заменили символы ? на значения без использования prepared statement, хотя это на деле это будет не так).

SELECT timestamp_rounded,
       min(air_temperature) AS min_air_temperature,
       max(air_temperature) AS max_air_temperature,
       avg(air_temperature) AS avg_air_temperature,
       median(air_temperature) AS median_air_temperature,
       varPop(air_temperature) AS var_air_temperature,
       min(air_humidity) AS min_air_humidity,
       max(air_humidity) AS max_air_humidity,
       avg(air_humidity) AS avg_air_humidity,
       median(air_humidity) AS median_air_humidity,
       varPop(air_humidity) AS var_air_humidity,
       min(wind_speed) AS min_wind_speed,
       max(wind_speed) AS max_wind_speed,
       avg(wind_speed) AS avg_wind_speed,
       median(wind_speed) AS median_wind_speed,
       varPop(wind_speed) AS var_wind_speed,
       min(illuminance) AS min_illuminance,
       max(illuminance) AS max_illuminance,
       avg(illuminance) AS avg_illuminance,
       median(illuminance) AS median_illuminance,
       varPop(illuminance) AS var_illuminance,
       (SELECT avg(rains) FROM (SELECT sum(is_raining) AS rains
                                FROM sensor_values
                                WHERE timestamp BETWEEN '2019-12-18 00:00:00' AND '2019-12-19 00:00:00'
                                    AND (altitude BETWEEN 150.0 AND 450.0)
                                GROUP BY latitude, longitude, altitude,
                                    toStartOfHour(timestamp))) AS avg_rains
FROM (SELECT timestamp, air_temperature, air_humidity, wind_speed, is_raining, illuminance
      FROM sensor_values
      WHERE timestamp BETWEEN '2019-12-18 00:00:00' AND '2019-12-19 00:00:00'
          AND (altitude BETWEEN 150.0 AND 450.0))
GROUP BY toStartOfHour(timestamp) AS timestamp_rounded
ORDER BY timestamp_rounded;

Можно повторно отметить, что ClickHouse принимает и хранит данные о времени с точностью до секунд, данные о миллисекундах отбрасываются.

Далее от ClickHouse возвращаются данные в формате TabSeparatedWithNamesAndTypes и на сервере, в том числе, с использованием официального ClickHouse JDBC, производятся те же операции, что были описаны в предыдущем подразделе.

После выполнения приведённого выше запроса клиенту вернётся ответ в формате JSON, приведённый в файле total-response.json. Подробнее о значении разных полей ответа можно почитать в API: в разделах про пользовательские типы данных и про метод метод POST /getAggregatedData. Поле click_house_query_execution_time представляет количество миллисекунд, в течение которых выполнялся запрос на получение агрегированных данных из ClickHouse.

Принципы работы генератора показаний датчиков

Так как нам не удалось найти реальные IOT-устройства для работы с реальными данными, было принято решение создать генератор "фейковых" показаний датчиков, расположенных в различных местностях. В нашей системе эти первичные данные периодически отправляются генератором, написанным на Node.js.

На данный момент в системе существует 15 различных датчиков, каждый из которых в порядке очереди отправляет раз в 2 секунды новые показания на основной сервер. Сделано это было для того, чтобы сымитировать передачу каждым из датчиков обновлённых данных раз в 30 секунд. Данные показаний датчиков передаются в следующем виде:

{
  "coordinates": {
    "latitude": 89.6345,
    "longitude": 3.8563234,
    "altitude": 12.9346
  },
  "values": {
    "air_temperature": 6.151137931377175,
    "air_humidity": 71.28986173922557,
    "wind_speed": 5.485463180331548,
    "is_raining": true,
    "illuminance": 5642.714283270595
  },
  "timestamp": "2019-12-19T20:58:00.435Z"
}

Кроме запланированной передачи данных в виде POST-запроса на основной сервер, описанного выше, существует публичное API для получения внеочередных данных с датчиков: /pollRandomSensor можно использовать для получения показаний с одного случайного датчика, а /pollAllSensors опрашивает все существующие датчики и снимает с них обновлённые данные.

В целях создания эффекта приближения динамики показаний датчиков к реальной жизни, изменения их показаний следуют следующей логике:

  • Для всех числовых показаний датчиков одна итерация изменений не может повлиять не предыдущее значение больше, чем на 20% от диапазона всех возможных значений в обе стороны, но может быть меньше, вплоть до 0%. Этот диапазон задаётся вручную. Уменьшается или увеличивается значение показания решает логика random_value >= 0.5. При невозможности произвести изменение значения в какую-либо из сторон, значение изменяется в другую сторону. Таким образом, мы можем гарантировать корректность значений и их более приближенную к реальности динамику в сравнении с простой рандомизацией.
  • Для всех булевых значений (в нашем случае, только для is_raining) существует очень низкая вероятность смены значения на противоположное - около 3%.

Для корректной работы генератора необходимо описать начальные значения датчиков и границы их изменений в этом файле. Пример такой конфигурации:

{
  "0": {
    "air_temperature": {
      "value": 15.819025107075214,
      "leftBoundary": 10,
      "rightBoundary": 20
    },
    "air_humidity": {
      "value": 68.60312366378226,
      "leftBoundary": 60,
      "rightBoundary": 80
    },
    "wind_speed": {
      "value": 3.834153577331329,
      "leftBoundary": 2,
      "rightBoundary": 7
    },
    "is_raining": {
      "value": true,
      "leftBoundary": false,
      "rightBoundary": true
    },
    "illuminance": {
      "value": 4780.0107244375295,
      "leftBoundary": 4700,
      "rightBoundary": 5100
    }
  }
  ...
}

После каждой произведённой модификации значений детчиков происходит перезапись файла с конфигурацией на диск с новыми значениями показаний и обновляются соответствующие данные в оперативной памяти. Необходимо это для того, чтобы минимизировать аномалии в изменениях данных при падении/перезагрузки облачной виртуальной машины.

Также, при неудачном POST запросе на основной сервер из-за внутренней проблемы на сервере (код ошибки 5xx) или из-за проблем в сетевом соединении 3 раза происходит переотправка данных для каждого датчика с интервалом в 30 секунд после последней попытки, чтобы не потерять сгенерированные данные.

Принципы работы клиента (мобильного iOS-приложения)

Пользователь вводит необходимые данные для отправки запроса на сервер:

  • Выбирает диапазон дат, за который нужны данные (параметр timestamp_interval).
  • Выбирает по каким отрезкам времени агрегировать данные дат (параметр aggregated_period ).
  • Опционально выбирает интервалы для местоположения датчиков - широта, долгота и высота. Параметры: latitude_intervals,longitude_intervals,altitude_intervals соответственно.

Скриншот с начальным экраном, после запуска приложения:

Скриншоты с вводом начальной даты:

Скриншоты с вводом конечной даты:

Скриншот с выбором времени группировки значений:

Скриншот с вводом опциональных интервалов для получения данных только от подходящих датчиков:

После ввода данных, приложение осуществляет запрос на сервер, и полученный ответ представляет в виде набора графиков по всем доступным значениям для каждой величины. Можно перейти на экран просмотра графиков по любой величине, а также просмотреть все значения на столбчатой диаграмме за выбранный период.

Скриншот мобильного приложения в процессе загрузки данных:

Скриншот с загруженными данными:

Скриншот столбчатой диаграммы по влажности воздуха:

Скриншот столбчатой диаграммы по влажности воздуха, после прокрутки вправо, чтобы посмотреть значения, не влезающие в экран:

Также можно вернуться назад на начальный экран, выбрать другие параметры для запроса и повторить запрос.

Эффективность использования ClickHouse в данной задаче

Генерация данных была запущена 15 декабря 2019 года в 21:09 по московскому времени.

На момент 19 декабря 2019 года 01:20 по московскому времени количество записей в единственной таблице sensor_values ClickHouse-кластера составляет 136515.

На этот же момент времени:

  • средний размер ежечасовой выгрузки в строках (среднее количество выгружаемых показаний датчиков) в ClickHouse, рассчитанный по данным из PostgreSQL-таблицы scheduled_pushing_data_to_clickhouse_log составляет 1796.25, причём в последних 72 выгрузках (из 76 всего имеющихся на данный момент) размер составляет ровно 1800 (что логично, поскольку каждые полминуты свои данные отправляют 15 датчиков, в одном часе 60 минут, а 15 * 2 * 60 = 1800);
  • средняя продолжительность выполнения запроса INSERT INTO при выгрузке данных в ClickHouse составляет 78.2236842105263158 миллисекунд с несмещённой оценкой дисперсии этой продолжительности 363.5892982456140351, минимум - 49 мс, а максимум - 149 мс, таким образом, вставка 1800 строк с 9 столбцами каждая (всего 1800 * 9 = 16200 ячеек) занимала доли секунды, не превышающие 0,15 секунды;
  • продолжительность выполнения запросов для получения агрегированных данных из ClickHouse варьируется в зависимости от параметров запроса - например, при получении агрегированных данных из ClickHouse, разделённых по отдельным датчикам (то есть с "need_split_by_coordinates": true в JSON-запросе от клиента к серверу), за один час с делением по минутам (то есть с "aggregated_period": "MINUTE" в JSON-запросе от клиента к серверу) либо за сутки с делением по часам (то есть с "aggregated_period": "HOUR" в JSON-запросе от клиента к серверу), время выполнения запроса, как правило, не превышает 2-3 секунд, а при повторных выполнениях того же запроса в течение некоторого времени после его первого выполнения может снизиться до нескольких сотен миллисекунд (то есть значительно менее 1 секунды), что, скорее всего, означает применение в ClickHouse технологий для кэширования результатов выполнения запросов, хотя официально на данный момент полноценная поддержка этого заявлена лишь в roadmap на 2020.

Выводы

Нам кажется, что ClickHouse даёт неплохие показатели на весьма большом количестве данных, даже при том, что для экономии средств нами в настоящее время используется ClickHouse-кластер с достаточно "слабыми" характеристиками (платформа Intel Cascade Lake типа burstable, 2 ядра vCPU с гарантированной долей 20 % vCPU и оперативной памятью 2 Гб) и аналогичная виртуальная машина для сервера (бэкенд-приложения) и генератора показаний датчиков (для экономии средств нами в настоящее время используется одна виртуальная машина на оба этих компонента), с такими же характеристиками, на Яндекс.Облаке.

Тем не менее, совсем "тяжёлые" запросы - например, получение агрегированных данных из ClickHouse, разделённых по отдельным датчикам (то есть с "need_split_by_coordinates": true в JSON-запросе от клиента к серверу), за несколько суток с делением по минутам (то есть с "aggregated_period": "MINUTE" в JSON-запросе от клиента к серверу) - в случае с нашим ClickHouse-кластером, могут вызывать такие ошибки, как, например, java.net.SocketTimeoutException: Read timed out на сервере и DB::Exception: Memory limit (total) exceeded в самом ClickHouse-кластере. Таким образом, хотя ClickHouse (как и другие OLAP-СУБД) и позволяет в некоторой мере сэкономить физические ресурсы, он не позволяет совсем отказаться от их "апгрейда" (например, при повышении требований к производительности и объёмов хранимых и обрабатываемых данных).