Войти

DTM-uploader

Дата актуализации: 29.01.2026.
Причины актуализации:
      – изменена информация об ответах во всех запросах;
      – добавлены значения service и source.

DTM-uploader – базовый сервис загрузки данных, который предоставляет возможность асинхронного приёма данных из сторонних источников с целью последующей загрузки их в Витрину  данных. Загрузка и обновление данных осуществляется в соответствии с заранее подготовленными Avro-схемами. Сервис загрузки данных реализуется компонентой ETL, предоставляющей REST API.

Для реализации передачи данных через сервис загрузки DTM-uploader, необходимо установить на Витрине данных компонент dtm_uploader (описание установки данного компонента приведено в статье – Установка витрины в конфигурации стандарт. Дополнительные компоненты).

Основные требования к передаваемым файлам

Загрузка данных в систему производится в виде сообщений, каждое из которых имеет структуру, представленную на рисунке 1:

Рисунок 1 - Структура загружаемых сообщений.jpg

Рисунок 1 – Структура загружаемых сообщений

Для успешной загрузки данные должны соответствовать следующим условиям:

  1.   Тело сообщения представляет собой файл Avro (Object Container File), который состоит из заголовка и блоков данных.

  2.   Заголовок файла содержит схему данных Avro.

  3.   Схема данных содержит следующие элементы:

  • имя;
  • тип “record”;
  • перечень полей.

  4.   Последним полем схемы должно быть указано служебное поле sys_op с типом данных avro.int.

  5.   Каждый блок данных содержит запись, представленную в бинарной кодировке.

  6.   Каждая запись содержит перечень полей и их значений.

  7.   Состав и порядок полей должны совпадать во всех следующих объектах:

  • в схеме данных заголовка файла Avro;
  • в наборе загружаемых записей;
  • во внешней таблице загрузки (поле sys_op должно отсутствовать);
  • в таблице-приемнике данных (поле sys_op должно отсутствовать).

Более подробно про формат данных Avro описано в источнике: https://avro.apache.org/docs/1.10.2/spec.html#Object+Container+Files.

  Примечание

 В загружаемой схеме данных Avro и записях Avro важны порядок и тип полей. Имена полей не сравниваются с именами полей внешней таблицы и таблицы-приемника

Приер ниже содержит схему данных Avro, используемую для загрузки данных о сотрудниках в таблицу staff. Для поля date_of_birth указан логический тип Avro, для поля middle_name — элемент union (поле является не обязательным для заполнения, поэтому маркер null выведен в отдельный параметр).

Пример схемы данных Avro:

{

  "name": "staff",

  "type": "record",

  "fields": [

    {

      "name": "id",

      "type": "long"

    },

    {

      "name": "firstname",

      "type": "string"

    },

    {

      "name": "lastname",

      "type": "string"

    },

    {

      "name": "middle_name",

      "type": [

        "null",

        "string"

      ]

    },

    {

      "name": "date_of_birth",

      "type": {

      "logicalType": "timestamp-micros",

      "type": "long"

       }  

    },

    {

      "name": "employee_position",

      "type": "string"

    },

    {

      "name": "department_category",

      "type": "long"

    },

    {

      "name": "sys_op",

      "type": "int"

    }

  ]

}

Пример ниже содержит набор записей о сотрудниках, загружаемых в таблицу staff. Для наглядности примера бинарные данные представлены в JSON-формате.

Пример записей Avro:

[

  {

    "id": 1000111,

    "firstname": "Елена",

    "lastname": "Фролова",

    "middle_name": "Андреевна",

    "date_of_birth": 4641084000000000,

    "employee_position": "Менеджер по подбору персонала",

    "department_category": 1,

    "description": "Сотрудники отдела кадров",

    "sys_op": 0

  },

  {

    "id": 1000005,

    "firstname": "Пётр",

    "lastname": "Платонов",

    "middle_name": "",

    "date_of_birth": 5639904000000000,

    "employee_position": "Руководитель отдела кадров",

    "department_category": 1,

    "description": "Сотрудники отдела кадров",

    "sys_op": 0

  }

]

Все операции с Prostore производятся в соответствии с порядком, который задаётся дельтой.

  Примечание

 Дельта - это целостная совокупность изменений в логической базе данных (включает в себя операции записи/изменения/удаления и имеет порядковый номер, уникальный в рамках логической базы данных). В рамках открытой дельты можно выполнить произвольное число операций записи. Нумерация дельт начинается с 0

Особенности работы с дельтами:

  • Для каждой логической базы одновременно может быть открыто не более одной дельты;
  • Не допускается загрузка разных изменений одного и того же набора данных в рамках одной дельты;
  • Изменения данных, производимые в рамках открытой дельты, изолированы от пользовательских запросов;
  • Если нужно вернуть состояние данных, которое предшествовало изменениям, выполненным в рамках открытой дельты, следует откатить дельту. Откат возможен только для открытой дельты, после закрытия дельты возврат к предыдущему состоянию недоступен.

Особенности реализации ETL

Функциональные особенности реализованного ETL включают в себя следующие пункты:

          1)     Генерация первичных ключей записей, передаваемых для загрузки, производится на стороне источника.

          2)     Каждая Avro-структура должна содержать данные только для одной таблицы Витрины.

          3)     В Avro-структурах данных источник заполняет тип операции sys_op:

  • 0 – для добавления новой или обновления существующей записи;
  • 1 – для удаления существующей записи (см. пример записей Avro).

          4)     ETL не выполняет преобразования данных, предназначенных для загрузки в Витрину данных.

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

          6)     Не должно быть двух версий одной записи в рамках одной дельты.

          7)     Нельзя менять порядок атрибутов в avro-схеме, поскольку данные при загрузке в БД распределяются в соответствии с тем перечнем, который был указан в avro-схеме.

          8)     Во все ответы, которые генерируются кодом ETL, добавляется строка "service": "DTM-uploader", которая нужна для того, чтобы понять, что это именно ответ от ETL, а не от промежуточных сервисов. Если такой строки нет, то возможны два варианта развития событий:

  • запрос не дошёл до ETL; 
  • ответ перехвачен сервисом, работающим между пользователем и ETL.

          9)     Во все ответы добавляется значение source, в котором передаётся подсистема-источник сообщения/события. Это значение нужно для разбора ошибок на стороне разработки ETL.

ЗАГРУЗКА / УДАЛЕНИЕ ДАННЫХ (newDelta, partOfDelta, data)

Сначала источник данных формирует и передаёт файлы по REST API, которые накапливаются Компонентом ETL. Далее данные загружаются и вставляются в Витрину данных через внешние таблицы, или удаляются. Статусы обработки каждой операции необходимо отслеживать через endpoint /status. Информация о каждом шаге процесса содержится в подразделах ниже.

НАЧАЛО ОПЕРАЦИИ ЗАГРУЗКИ / УДАЛЕНИЯ СОГЛАСОВАННЫХ ДАННЫХ (newDelta, partOfDelta, data)

Endpoint – newDelta

Под endpoint’ом /newDelta регистрируется новая порция данных. Для того чтобы начать работу с данными, источнику данных необходимо сгенерировать UUID (идентификатор для новой порции данных) и вставить его в запрос c endpoint’ом /newDelta.

 Согласованные данные – это данные для нескольких таблиц, которые должны попасть в Витрину данных за одну операцию вставки (то есть в одной дельте), а значит будут доступны для потребителя одномоментно (такой способ загрузки актуален, когда необходимо обновить данные)

Пример набора данных, который будет загружен или удален в рамках одной дельты представлен ниже:

{

     "requestId": "6B29FC40-CA47-1067-B31D-00DD010662DA",

     "dataSetName": ["product", "stock"]

}

где:

  • <requestId> — идентификатор порции изменений (дельты);
  • <dataSetName> — массив имен набора данных (product и stock - имена таблиц).

Запрос с endpoint’ом /newDelta будет иметь вид:

curl -X POST "http://localhost:9999/newDelta" -H 'Content-Type: application/json' -d '{"requestId": "6B29FC40-CA47-1067-B31D-00DD010662DA", "dataSetName": ["product", "stock"]}'

где:

  • <requestId> — идентификатор порции изменений (дельты), который был ранее сгенерирован;
  • <dataSetName> — имя набора данных (имена таблиц).

  Примечание

 Данные для обновления/вставки и для удаления должны быть отдельно зарегистрированы в newDelta с разными requestId. Если требуется обновить/вставить данные, то сначала нужно зарегистрировать порцию данных через newDelta с requestId, затем прислать данные с зарегистрированным requestId через endpoint /partOfDelta. Допускается передача как в одном файле одним запросом, так и в двух файлах двумя запросами - это не имеет значения. Главное условие - нужно отправлять только данные на обновление и вставку, и они должны быть уникальны по первичному ключу. Также в них не должно быть одинаковых записей с одним первичным ключом. Пример: отправлены два обновления одной записи (с одинаковым первичным ключом). В этом случае в хранилище Prostore попадет только одна запись, и неизвестно какая, поэтому дублей по первичному ключу отправленных с одним requestId быть не должно!

Чтобы проверить статус выполнения запроса, необходимо направить запрос с endpoint’ом /status.

Если обработка запроса завершится успешно, то Витрина данных вернёт JSON-ответ, содержащий статус-сообщение об успешной операции:

{

    "requestId": "6B29FC40-CA47-1067-B31D-00DD010662DA",

    "dataSets": [

     "product",

     "stock"

    ],

    "status": "WAIT_DATA",

    "statusMessage": "Сервис ожидает данные",

    "inDeltaFlag": true,

    "errors": []

    "service": "DTM-uploader",

    "source": "DTM_UPLOADER",

    "statusCode": "SUCCESS"

}

где:

  • <requestId> — идентификатор порции изменений (дельты);
  • <status> — код возвращаемого статуса дельты (SUCCESS – запрос выполнен успешно);
  • <statusMessage> — код возвращаемого статуса дельты (WAIT_DATA –дельта готова принять данные);
  • <service> — значение всегда "DTM-uploader";
  • <source> — источник сообщения (DTM_UPLOADER – сообщение сгенерировано ETL);
  • <statusCode> — код выполнения операции получения статуса дельты (SUCCESS – код успешного выполнения операции, ERROR – ошибка получения статуса дельты).  

Пример JSON-ответа на проверку статуса, завершившегося ошибкой, приведен ниже:

{

     "requestId": "6B29FC40-CA47-1067-B31D-00DD010662DA",

     "statusCode": "REQUESTID_ALREADY_EXIST",

     "service": "DTM-uploader",

     "source": "DTM_UPLOADER",

     "statusMessage": "Дельта с requestId = 6B29FC40-CA47-1067-B31D-00DD010662DA уже существует. Статус загрузки дельты: SUCCESS"

}

где:

  • <requestId> — идентификатор порции изменений (дельты) который необходимо проверить;
  • <statusCode> — код возвращаемого статуса (REQUESTID_ALREADY_EXIST – идентификатор уже существует в БД);
  • <service> — значение всегда "DTM-uploader";
  • <source> — подсистема-источник сообщения/события;
  • <statusMessage> — сообщение с описанием кода ошибки.

Описание возвращаемых кодов:

  • ERROR – внутренняя ошибка (http code 500);
  • SUCCESS – успешное выполнение (http code 200);
  • REQUESTID_ALREADY_EXIST – requestId уже зарегистрирован (http code 400);
  • PROCESSING – идет обработка данных (http code 400).      

  Примечание

 Для удаления записей необходимо зарегистрировать новую дельту во endpoint’е /newDelta с новым requestId, и по зарегистрированному requestId должны быть присланы только данные на удаление

ЗАГРУЗКА/УДАЛЕНИЕ СОГЛАСОВАННЫХ ДАННЫХ

 Endpoint – partOfDelta

На данном шаге выполняется накапливание порции данных, создаётся вставка в Витрину данных по флагу isLastChunk.

  Примечание

 Если в процессе загрузки вызван метод newDelta, то текущая загрузка будет прервана и порция не попадет в Витрину данных

Чтобы отправить последнюю порцию данных для таблицы product, необходимо направить запрос с endpoint’ом /partOfDelta с указанием dataSetName=product и isLastChunk=true (что означает что данная порция данных - последняя). Обработка и загрузка данных не начнётся, пока не будет направлен такой же запрос, но уже по таблице stock: dataSetName=stock, isLastChunk=true.

Пример запроса на загрузку данных под ранее созданный набор данных:

curl -X POST "http://localhost:9999/partOfDelta" -F upload=@"./product.avro" -F dataSetName=product -F chunkNumber=0 -F isLastChunk=false -F requestId=a6212a7d-4526-4e2d-89a7-9828f380c91d

где:

  • <upload> — загружаемый avro-файл (пример avro-файла с данными представлен в разделе 2);
  • <dataSetName> — имя набора данных (имя таблицы);
  • <chunkNumber> — номер порции dataSet в рамках дельты;
  • <isLastChunk> — флаг последней порции dataSet;
  • <requestId> — идентификатор порции изменений (дельты).

В результате успешной загрузки при проверке статуса (пример вызова /statusrequestId на запрос Витрина данных вернёт JSON-ответ, содержащий статус-сообщение об успешной операции:

{

    "requestId": "f3947645-88c8-4044-bd8b-de273f8a8461",

    "statusCode": "SUCCESS",

    "statusMessage": "Порция получена.",

    "service": "DTM-uploader",

    "source": "DTM_UPLOADER"

}

где:

  • <requestId> — идентификатор порции изменений (дельты);
  • <statusCode> — статус код результата запроса (SUCCESS - запрос выполнен успешно);
  • <statusMessage> — описание статусного сообщения;
  • <service> — значение всегда "DTM-uploader";

  • <source> — подсистема-источник сообщения/события.

Если загрузка прервалась ошибкой, то при проверке статуса (пример вызова /statusrequestId на запрос Витрина данных вернёт JSON-ответ с описанием ошибки:

{

    "requestId": "aef2f195-b037-4aaa-b171-f2746511e7e2",

    "dataSets": [

        "stock"

    ],

    "inDeltaFlag": true,

    "status": "ERROR",

    "statusMessage": "Произошла ошибка"

    "errors": [

        {

            "dataSet": "stock",

            "errorType": "INSERT",

            "message": "Ошибка вставки в таблицы: stock"

        }

    ],

    "service": "DTM-uploader",

    "source": "PROSTORE",

    "statusCode": "SUCCESS"

}

где:

  • <requestId> — идентификатор порции изменений (дельты);
  • <dataSets> — массив имен набора данных (имен таблиц где была допущена ошибка);
  • <inDeltaFlag> = true — загрузка согласованных данных производилась через endpoint /partOfDetla;
  • <status> — статус код результата запроса (ERROR – внутренняя ошибка);
  • <statusMessage> — описание статусного сообщения;
  • <errors> — массив, ошибки загрузки или парсинга входящих данных;
  • <dataSet> — название таблицы где допущена ошибка;
  • <errorType> — тип ошибки;
  • <message> — описание ошибки;
  • <service> — значение всегда "DTM-uploader";

  • <source> — подсистема-источник сообщения/события;
  • <statusCode> — код выполнения операции получения статуса дельты (SUCCESS – код успешного выполнения операции, ERROR – ошибка получения статуса дельты).

Описание возвращаемых кодов:

  • NOT_FOUND - данные не найдены, либо были утеряны в результате остановки сервиса (http code 400);
  • REQUESTID_ALREADY_EXIST – requestId уже зарегистрирован – (http code 400);
  • PROCESSING – идет обработка данных (http code 400);
  • WRONG_ENDPOINT – requestId зарегистрирован для другого endpoint’а (http code 400);
  • EMPTY_ATTACHMENT – нет файла вложения (http code 400);
  • UNREGISTERED_DATASETNAME незарегистрированный набор данных (http code 400);
  • ERROR - внутренняя ошибка (http code 500);
  • SUCCESS - успешное выполнение (http code 200).

  Примечание

 Возможна ситуация, когда после падения ETL приходит запрос с requestId, который был до падения, в данном случае Витрина данных возвращает ошибку со статусом NOT_FOUND. Необходимо снова направить запрос по endpoint’у /newDelta с новым requestId и начать процесс загрузки заново

ЗАГРУЗКА/УДАЛЕНИЕ НЕСОГЛАСОВАННЫХ ДАННЫХ

Endpoint – data

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

 Несогласованные данные – могут быть вставлены в разных дельтах и будут доступны потребителю постепенно по мере загрузки. Этот способ подходит для первоначальной загрузки, когда еще нет потребителей

Вставка в Витрину данных выполнится после накопления порции или по флагу isLast, который используется для последней порции данных. Флаг isLast подаёт сигнал для завершения формирования дельты, для того чтобы выполнить вставку накопленных данных и закрыть транзакцию.

Пример запроса:

curl -X POST "http://localhost:9999/data" -F upload=@"./product.avro" -F dataSetName=product -F isLast=false -F requestId=a6212a7d-4526-4e2d-89a7-9828f380c91d

где:

  • <upload> — загружаемый avro-файл (пример avro-файла с данными представлен в разделе 2);
  • <dataSetName> — имя набора данных (имя таблицы);
  • <isLast> — флаг последней порции данных (сигнал для завершения формирования дельты, для того чтобы выполнить вставку накопленных данных и закрыть транзакцию.);
  • <requestId> — идентификатор порции изменений (дельты).

В результате успешной операции при проверке статуса запроса по endpoint’у /status Витрина данных вернёт JSON-ответ, содержащий статус-сообщение:

{

     "requestId": "6B29FC40-CA47-1067-B31D-00DD010662DA",

     "statusCode": "SUCCESS",

     "statusMessage": "Порция получена.",

     "service": "DTM-uploader",

     "source": "DTM_UPLOADER"

}

где:

  • <requestId> — идентификатор порции изменений (дельты);
  • <statusCode> — статус код результата запроса (SUCCESS - запрос выполнен успешно);
  • <statusMessage> — описание статусного сообщения;
  • <service> — значение всегда "DTM-uploader";

  • <source> — подсистема-источник сообщения/события.

Если загрузка прервалась ошибкой, то при проверке (пример вызова /statusrequestId Витрина данных вернёт JSON-ответ с описанием ошибки:

{

     "requestId": "6B29FC40-CA47-1067-B31D-00DD010662DA",

     "statusCode": "NOT_FOUND",

     "statusMessage": "Не найдена дельта с requestId = 6B29FC40-CA47-1067-B31D-00DD010662DA",

     "service": "DTM-uploader",

     "source": "DTM_UPLOADER"

}

где:

  • <requestId> — идентификатор порции изменений (дельты);
  • <statusCode> — статус код результата запроса (NOT_FOUND - данные по requestId были утеряны в результате остановки сервиса, необходимо зарегистрировать новую дельту и снова загрузить данные);
  • <statusMessage> — описание статусного сообщения;
  • <service> — значение всегда "DTM-uploader";
  • <source> — подсистема-источник сообщения/события.

ПРОВЕРКА СТАТУСНОЙ ИНФОРМАЦИИ ПО ЗАГРУЗКЕ / УДАЛЕНИЮ ДАННЫХ

Endpoint – status

Endpoint /status предназначен для проверки статусной информации из сервисных таблиц по requestId.

Пример запроса:

curl -X GET "http://localhost:9999/status" -d "requestId=13f2475e-f3dc-4c9e-b2f6-3a98320261f1"

где:

  • <requestId> — UUID идентификатор порции изменений (дельты).

Пример ответа на такой запрос представлен ниже:

{

    "requestId": "13f2475e-f3dc-4c9e-b2f6-3a98320261f1",

    "inDeltaFlag": false,

    "dataSets": [

        "stock"

    ],

    "status": "ERROR",

    "statusMessage": "Произошла ошибка",

    "errors": [

        {

            "dataSet": "stock",

            "errorType": "PARCING",

            "message": "Неверно указан тип поля count_pieces: LONG. Ожидается: INTEGER"

        },

        {

            "dataSet": "stock",

            "errorType": "PARCING",

            "message": "Неверно указан тип поля product_id: LONG. Ожидается: INTEGER"

        }

    ],

    "prostoreInfo":"{"get_delta_hot": [{"delta_num":68,"cn_from":164,"cn_max":165,"is_rolling_back":false,"write_op_finished":"[{"tableName":"users","cnList":[{"cn":165,"status":0,"rowsAffected":2}]}]"}], "get_delta_ok": [{"delta_num":68,"delta_date":"2025-08-14 06:37:07.545","cn_from":164}]}",

    "service":"DTM-uploader",

    "source":"DTM_UPLOADER",

    "statusCode":"SUCCESS"

}

где:

  • <requestId> — UUID идентификатор порции изменений (дельты);
  • <inDeltaFlag = false > — загрузка несогласованных данных производилась через endpoint /data;
  • <dataSets> — массив имен набора данных (имен таблиц где была допущена ошибка);
  • <status> — статус код результата запроса дельты (SUCCESS);

  Важно!

 Если после загрузки всех данных дельты получен статус ERROR, то в некоторых случаях, в частности, когда ошибки вызваны проблемами на сетевом оборудовании, имеет смысл повторить отправку всех данных, входящих в дельту, предварительно сгенерировав новый UUID идентификатор порции изменений requestId


  • <statusMessage> — описание статусного сообщения;
  • <prostoreInfo> - информация о сохранённой дельте Prostore. Если дельта ещё не сохранялась, поле будет содержать значение «no data». Если дельта сохранилась, то выведется результат выполнения к Prostore запросов get_delta_hot (непосредственно перед сохранением дельты) и get_delta_ok -–сразу после сохранения дельты;

  • <service> — значение всегда "DTM-uploader";
  • <source> — подсистема-источник сообщения/события;
  • <statusCode> — код выполнения операции получения статуса дельты (SUCCESS – код успешного выполнения операции, ERROR – ошибка получения статуса дельты);
  • <errors> — массив, ошибки загрузки или парсинга входящих данных;
  • <dataSet> — название таблицы, где допущена ошибка;
  • <errorType> — тип ошибки;
  • <message> — описание ошибки.

Таким образом, для любых операций с данными посредством ETL, необходимо:

          1)     Подготовить Avro-файлы с данными;
          2)     Сформировать произвольный UUID, в рамках которого зарегистрируется новая порция данных. Направить запрос по endpoint’ту /newDelta (указав таблицы, в рамках которых планируются изменения);
          3)     Сформировать и направить запрос(ы) в рамках которых будут накапливаться данные на загрузку в рамках одной операции.

РАБОТА С ВЛОЖЕНИЯМИ ЧЕРЕЗ S3

Загрузка данных в хранилище

Endpoint – uploadAttachment

Перед загрузкой источнику данных необходимо сгенерировать UUID (идентификатор для новой порции данных) и вставить его в запрос с endpoint’ом /uploadAttachment. При совпадении имен вложений в хранилище, вложение перезаписывается.

Пример запроса на загрузку вложения в хранилище представлен ниже:

curl -X POST "http://localhost:9090/uploadAttachment" -F upload=@"document.pdf" -F requestId=13f2475e-f3dc-4c9e-b2f6-3a98320261f1 -F name=Doc_1

где:

  • <upload> — путь до загружаемого файла-вложения;
  • <requestId> — UUID идентификатор запроса;
  • <name> уникальное имя вложения.

Метод синхронный, результат загрузки файла возвращается в ответе на запрос.

После успешной загрузки при проверке по endpoint’у /status Витрина данных вернёт JSON-ответ, содержащий статус-сообщение:

{

    "requestId": "13f2475e-f3dc-4c9e-b2f6-3a98320261f1",

    "statusCode": "UPDATED",

    "statusMessage": "Файл Doc_1 успешно обновлен.",

    "service":"DTM-uploader",

    "source":"S3"

}

где:

  • <requestId> — UUID идентификатор запроса;
  • <statusCode> — статус код результата запроса (SUCCESS - запрос выполнен успешно; UPDATED – данные обновлены);
  • <statusMessage> — описание статусного сообщения;
  • <service> — значение всегда "DTM-uploader";
  • <source> — подсистема-источник сообщения/события.

Пример неуспешной загрузки представлен ниже:

{

    "requestId": "13f2475e-f3dc-4c9e-b2f6-3a98320261f1",

    "statusCode": "ERROR",

    "statusMessage": "Произошла ошибка",

    "service":"DTM-uploader",

    "source":"S3"

}

где:

  • <requestId> — UUID идентификатор запроса;
  • <statusCode> — статус код результата запроса (ERROR – запрос завершился ошибкой);
  • <statusMessage> — описание статусного сообщения;
  • <service> — значение всегда "DTM-uploader";
  • <source> — подсистема-источник сообщения/события.

Описание возвращаемых кодов:

  • UPDATED – данные обновлены (http code 200);
  • EMPTY_ATTACHMENT – нет файла вложения (http code 400);
  • ERROR – внутренняя ошибка (http code 500);
  • SUCCESS – успешное выполнение (http code 200).

Удаление данных из хранилища 

Endpoint – deleteAttachment


Для того чтобы удалить вложения из хранилища S3 необходимо направить следующий запрос:

curl -X DELETE "http://localhost:9090/deleteAttachment?name=aef2f195-b037-4aaa-b171-f2746511e7e2"

где:

  • <name> — UUID идентификатор запроса.

Метод синхронный, результат удаления файла возвращается в ответе на запрос.

В результате успешного удаления Витрина данных вернёт JSON-ответ, содержащий статус-сообщение.

{

    "requestId": "13f2475e-f3dc-4c9e-b2f6-3a98320261f1",

    "statusCode": "SUCCESS",

    "statusMessage": "Файл Doc_1 успешно удален.",

    "service":"DTM-uploader",

    "source":"S3"

}

где:

  • <requestId> — UUID идентификатор запроса;
  • <statusCode> — статус код результата запроса (SUCCESS - запрос выполнен успешно);
  • <statusMessage> — описание статусного сообщения;
  • <service> — значение всегда "DTM-uploader";

  • <source> — подсистема-источник сообщения/события.

Если удаление завершилось ошибкой, то Витрина данных вернёт JSON-ответ c кодом ошибки:

{

    "requestId": "13f2475e-f3dc-4c9e-b2f6-3a98320261f1",

    "statusCode": "NOT_FOUND",

    "statusMessage": "Файл Doc_1 не найден.",

    "service":"DTM-uploader",

    "source":"S3"

}

где:

  • <requestId> — UUID идентификатор запроса;
  • <statusCode> — статус код результата запроса (NOT_FOUND);
  • <statusMessage> — описание статусного сообщения;
  • <service> — значение всегда "DTM-uploader";

  • <source> — подсистема-источник сообщения/события.

МАППИНГ ДАННЫХ 

Endpoint – generateMapping

Endpoint предназначен для первичной настройки, а также перенастройки сервиса в случае изменения модели данных Витрины. Файл маппинга генерируется источником данных в формате Kotlin-script. В файле описана модель данных Витрины данных в виде структуры объектов Kotlin (table, column). Объекты table описывают таблицы, каждый из них содержит имя таблицы и список колонок в том порядке, в котором они созданы в Витрине. Объекты column описывают колонки, каждый из них содержит имя колонки, тип данных, признак обязательности (nullable), признак первичного ключа.

Файл используется сервисом для описания модели данных и валидации входящих данных. Выполняются следующие проверки:

  • проверяется соответствие состава полей входящей avro-структуры составу полей, описанных в файле маппинга;
  • проверяется соответствие порядка полей входящей avro-структуры порядку полей, описанных в файле маппинга;
  • проверяется соответствие типов данных полей входящей avro-структуры типам полей, описанных в файле маппинга. Для полей с установленным признаком обязательности (nullable = false) выполняется проверка на null.

При вызове endpoint’а /generateMapping сервис генерирует файл на основе информации о модели, полученной из развернутой Витрины данных. Файл складывается сервисом на диск, а также возвращается в ответе на вызов.

Пример запроса на генерацию маппинга представлен ниже:

curl -X GET "http://localhost:9090/generateMapping"

Результат Витрина данных вернёт в формате Kotlin-script:

import ru.supercode.mapping.common.ColumnType.*

import ru.supercode.mapping.mapper.dsl.mappingAvro

mappingAvro {

    table("product") {

        column("id", INTEGER) { nullable = false; primary = true; }

        column("name", STRING) { nullable = false; }

    }

    table("stock") {

        column("product_id", INTEGER) { nullable = false; primary = true; }

        column("count_pieces", INTEGER) { nullable = false; }

    }

}

ВАЛИДАЦИЯ ДАННЫХ

Валидация порции данных производится в момент обработки и вставки.

  Примечание

 Помимо валидации данных осуществляется валидация параметров запроса. Во всех endpoint’ах requestId должен быть в формате UUID

В случае ошибок при валидации результат будет возвращен при вызове endpoint’а /status.

Ошибки, возникающие в процессе обработки endpoint’а /newDelta:

  • отклоняются запросы, которые получены в момент обработки порции данных (statusCode: PROCESSED);
  • если пустой параметр dataSetName;
  • прислан запрос с уже зарегистрированным requestId и statusCode данного requestId не равен NOT_FOUND или WAIT_DATA.

Ошибки, возникающие в процессе обработки endpoint’а /partOfDelta:

  • прислан запрос с незарегистрированным requestId;
  • прислан запрос с уже зарегистрированным requestId и statusCode данного requestId не равен NOT_FOUND или WAIT_DATA;
  • прислан запрос с requestId зарегистрированным для endpoint'а /data;
  • прислан запрос с параметром dataSetName, который не был зарегистрирован в endpoint’е /newDelta;
  • нет файла вложения в параметре upload.

Ошибки, возникающие в процессе обработки endpoint’а /data:

  • отклоняются запросы, которые получены в момент обработки порции данных (statusCode: PROCESSED);
  • прислан запрос с незарегистрированным requestId;
  • прислан запрос с уже зарегистрированным requestId и statusCode данного requestId не равен NOT_FOUND или WAIT_DATA;
  • прислан запрос с requestId зарегистрированным для endpoint’а /partOfDelta;
  • нет файла вложения в параметре upload.

Ошибка, возникающая в процессе обработки endpoint’а /uploadAttachment:

  • Нет файла вложения в параметре upload.

Ошибка, возникающая в процессе обработки endpoint’а /generateMapping:

  • Не созданы логические таблицы в схеме.
Авторизуйтесь, чтобы оставить комментарий к статье