Диаграммы последовательностей работы коннектора Meta#

Данный документ содержит диаграммы последовательностей (PlantUML), описывающие алгоритмы работы коннектора Meta (IDMC Connector).

Тестирование подключения (test)#

@startuml
title Диаграмма последовательности: Тестирование подключения

participant "IDM Framework" as Framework
participant "MetaConnector" as Connector
participant "TopicReader" as Reader
participant "KafkaConfiguration" as KafkaConfig
participant "Consumer" as Consumer
participant "Kafka Broker" as Kafka

Framework -> Connector: test()
activate Connector

Connector -> Reader: checkPartitions()
activate Reader

Reader -> KafkaConfig: getPartitions(config)
activate KafkaConfig
KafkaConfig -> KafkaConfig: AdminClient.describeTopics()
KafkaConfig --> Reader: List<TopicPartition>
deactivate KafkaConfig

Reader -> Consumer: endOffsets(partitions, duration)
activate Consumer
Consumer -> Kafka: запрос end offsets
Kafka --> Consumer: Map<TopicPartition, Long>
deactivate Consumer

alt успех
    Reader --> Connector: OK
    Connector --> Framework: OK
else исключение
    Reader --> Connector: Exception
    Connector -> Connector: логировать ошибку
    Connector --> Framework: ConnectorException
end

deactivate Reader
deactivate Connector
@enduml

Выполнение поиска (executeQuery)#

@startuml
title Диаграмма последовательности: Выполнение поиска (executeQuery)

participant "IDM Framework" as Framework
participant "MetaConnector" as Connector
participant "ExecuteQueryHandler" as ExecH
participant "TopicReader" as Reader
participant "TopicSender" as Sender
participant "Consumer" as Consumer
participant "Kafka Broker" as Kafka
participant "JsonSchema" as JsonSchema
participant "ResultsHandler" as Results

Framework -> Connector: executeQuery(objectClass, filter, resultsHandler, options)
activate Connector
Connector -> ExecH: getObjects(objectClass, resultsHandler, filter)
activate ExecH

ExecH -> Reader: getConsumerRecords(filter)
activate Reader

Reader -> Reader: getPartitions()
Reader -> Consumer: assign(partitions)

alt checkOffsetValue = true
    Reader -> Consumer: endOffsets(partitions)
    Reader -> Consumer: position(partition)
    alt offset > end
        Reader -> Consumer: seekToEnd(partitions)
        Reader -> Consumer: commitSync()
    end
end

Reader -> Consumer: poll(duration)
Consumer -> Kafka: fetch request
Kafka --> Consumer: ConsumerRecords
Reader --> ExecH: List<ConsumerRecord>
deactivate Reader

loop для каждого ConsumerRecord
    ExecH -> ExecH: processRecord(record)
    
    alt record.value != null
        ExecH -> JsonSchema: validate(message)
        JsonSchema --> ExecH: Set<ValidationMessage>
        
        alt report.isEmpty()
            ExecH -> ExecH: validateDate()
            ExecH -> ExecH: validateVersions()
            ExecH -> ExecH: parseJson()
            ExecH -> ExecH: buildConnectorObject()
            ExecH -> Results: handle(connectorObject)
        else ошибки валидации
            ExecH -> ExecH: buildInvalidConnectorObject()
            ExecH -> Results: handle(invalidObject)
            ExecH -> Sender: sendMessage(error)
        end
    else record.value == null
        ExecH -> Sender: sendMessage("Record value is null")
    end
end

ExecH --> Connector: завершено
deactivate ExecH
deactivate Connector
@enduml

Валидация сообщения и отправка результата#

@startuml
title Диаграмма последовательности: Валидация сообщения и отправка результата

participant "ExecuteQueryHandler" as ExecH
participant "JsonSchema" as JsonSchema
participant "TopicSender" as Sender
participant "Producer" as Producer
participant "Kafka Broker" as Kafka

ExecH -> JsonSchema: validate(kafkaMessage)
JsonSchema --> ExecH: Set<ValidationMessage>

alt report.isEmpty()
    ExecH -> ExecH: validateDate(kafkaMessage, config)
    alt дата не валидна
        ExecH -> ExecH: buildInvalidConnectorObject()
        ExecH -> Sender: sendMessage(errorMessage)
    end
    
    ExecH -> ExecH: validateVersions(kafkaMessage)
    alt версии не валидны
        ExecH -> ExecH: buildInvalidConnectorObject()
        ExecH -> Sender: sendMessage(errorMessage)
    end
    
    ExecH -> ExecH: parseJson()
    ExecH -> ExecH: buildConnectorObject()
else report не пуст
    ExecH -> ExecH: создать Map ошибок
    ExecH -> Sender: sendMessage(Map)
    ExecH -> ExecH: buildInvalidConnectorObject()
end

== Отправка сообщения ==
ExecH -> Sender: send(message, topic)
activate Sender

alt producerNameOfTopic не пуст
    Sender -> Producer: createProducer()
    Sender -> Producer: send(ProducerRecord)
    Producer -> Kafka: отправить сообщение
    Kafka --> Producer: acknowledgment
    Sender --> ExecH: OK
else producerNameOfTopic пуст
    Sender -> Sender: логировать ошибку
end
deactivate Sender
@enduml

Выполнение скрипта на ресурсе (runScriptOnResource)#

@startuml
title Диаграмма последовательности: Выполнение скрипта на ресурсе

participant "IDM Framework" as Framework
participant "MetaConnector" as Connector
participant "ScriptHandler" as ScriptH
participant "KafkaExecutorHandler" as ExecutorH
participant "JsonSchema" as JsonSchema
participant "TopicSender" as Sender
participant "Consumer" as Consumer
participant "Kafka Broker" as Kafka

Framework -> Connector: runScriptOnResource(scriptContext, options)
activate Connector
Connector -> ScriptH: runScriptOnResource(scriptContext, options)
activate ScriptH

ScriptH -> ScriptH: проверить scriptLanguage

alt language == "KAFKA_EXECUTOR"
    ScriptH -> ExecutorH: new KafkaExecutorHandler(config, sender)
    activate ExecutorH
    
    ExecutorH -> JsonSchema: загрузить KafkaExecutor.json schema
    ExecutorH -> ExecutorH: validateScriptCode(script)
    ExecutorH -> JsonSchema: validate(script)
    
    alt валидация неудачна
        JsonSchema --> ExecutorH: ValidationMessage[]
        ExecutorH --> Framework: ConfigurationException
        deactivate ExecutorH
        deactivate ScriptH
        deactivate Connector
    end
    
    ExecutorH -> ExecutorH: распарсить commands (ArrayNode)
    
    loop для каждой команды
        ExecutorH -> ExecutorH: parseAction(command)
        
        alt command == "SEND_MESSAGE"
            ExecutorH -> ExecutorH: sendMessage()
            ExecutorH -> Sender: send(message, producerNameOfTopic)
            Sender -> Kafka: отправить сообщение
        else command == "SEND_VERIFICATION_RESULT"
            ExecutorH -> ExecutorH: sendVerificationResult()
            ExecutorH -> JsonSchema: validate(result_message)
            ExecutorH -> Sender: send(message, verificationResultTopic)
            Sender -> Kafka: отправить сообщение
        else command == "COMMIT"
            ExecutorH -> ExecutorH: commit()
            ExecutorH -> Consumer: commitSync()
            Consumer -> Kafka: commit offset
        else command == "GET_OFFSET"
            ExecutorH -> ExecutorH: getOffset()
            ExecutorH -> Consumer: position(partition)
            Consumer --> ExecutorH: offset
        else command == "SEEK"
            ExecutorH -> ExecutorH: seek()
            ExecutorH -> Consumer: seek(partition, offset)
            Consumer -> Kafka: установка offset
        end
    end
    
    ExecutorH --> ScriptH: result
    ExecutorH -> ExecutorH: dispose()
    ExecutorH -> Consumer: close()
    deactivate ExecutorH
else
    ScriptH --> Framework: ConnectorException("Script language not supported")
    deactivate ScriptH
    deactivate Connector
end

ScriptH --> Framework: result
deactivate ScriptH
deactivate Connector
@enduml

Отправка сообщения верификации (SEND_VERIFICATION_RESULT)#

@startuml
title Диаграмма последовательности: Отправка сообщения верификации

participant "KafkaExecutorHandler" as ExecutorH
participant "JsonSchema" as JsonSchema
participant "TopicSender" as Sender
participant "Producer" as Producer
participant "Kafka Broker" as Kafka

ExecutorH -> ExecutorH: parseAction("SEND_VERIFICATION_RESULT")
ExecutorH -> ExecutorH: sendVerificationResult(objectNode)

ExecutorH -> ExecutorH: получить result_verification_message
ExecutorH -> JsonSchema: validate(result_message)
JsonSchema --> ExecutorH: Set<ValidationMessage>

alt result не пуст (ошибки)
    ExecutorH --> Framework: ConfigurationException
end

ExecutorH -> Sender: send(result_message, verificationResultTopic)
activate Sender

Sender -> Producer: send(ProducerRecord)
Producer -> Kafka: отправить сообщение
Kafka --> Producer: acknowledgment

Sender --> ExecutorH: OK
deactivate Sender
@enduml

Чтение сообщений из Kafka (TopicReader.getConsumerRecords)#

@startuml
title Диаграмма последовательности: Чтение сообщений из Kafka

participant "ExecuteQueryHandler" as ExecH
participant "TopicReader" as Reader
participant "Consumer" as Consumer
participant "Kafka Broker" as Kafka

ExecH -> Reader: getConsumerRecords(filter)
activate Reader

Reader -> Reader: getPartitions()
Reader --> Reader: List<TopicPartition>

alt filter != null
    Reader -> Reader: applyFilter(partitions, filter)
end

loop для каждой partition
    Reader -> Reader: createConsumer(partition)
    Reader -> Consumer: assign(partitions)
    
    alt checkOffsetValue = true
        Reader -> Consumer: endOffsets(partitions)
        Consumer --> Reader: Map<TopicPartition, Long>
        
        Reader -> Consumer: position(partition)
        Consumer --> Reader: offset
        
        alt offset > end
            Reader -> Consumer: seekToEnd(partitions)
            Reader -> Consumer: commitSync()
            Consumer -> Kafka: commit offset
            Kafka --> Consumer: OK
        end
    end
    
    Reader -> Consumer: poll(Duration)
    Consumer -> Kafka: fetch request
    Kafka --> Consumer: ConsumerRecords
    
    loop для каждого ConsumerRecord
        Reader -> Reader: добавить в result
    end
    
    alt autoCommit = true
        Reader -> Consumer: commitSync()
        Consumer -> Kafka: commit offset
        Kafka --> Consumer: OK
    end
end

Reader --> ExecH: List<ConsumerRecord>
deactivate Reader
@enduml

Отправка сообщения в Kafka (TopicSender.send)#

@startuml
title Диаграмма последовательности: Отправка сообщения в Kafka

participant "ExecuteQueryHandler" as ExecH
participant "TopicSender" as Sender
participant "Producer" as Producer
participant "Kafka Broker" as Kafka

ExecH -> Sender: send(message, topic)
activate Sender

Sender -> Sender: UUID.randomUUID() (key)
Sender -> Producer: send(ProducerRecord)
Producer -> Kafka: отправить сообщение
Kafka --> Producer: acknowledgment

alt успех
    Producer --> Sender: OK
    Sender --> ExecH: OK
else InterruptedException
    Sender -> Sender: логировать ошибку
    Sender --> ExecH: ConnectorException
else другое исключение
    Sender -> Sender: логировать ошибку
    Sender --> ExecH: ConnectorException
end

deactivate Sender
@enduml