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

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

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

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

participant "IDM Framework" as Framework
participant "KafkaConnector" as Connector
participant "TopicReader" as Reader
participant "KafkaConfigManager" as ConfigMgr
participant "AdminClient" as Admin
participant "Kafka Broker" as Kafka

Framework -> Connector: test()
activate Connector
Connector -> Reader: testConnection()
activate Reader

Reader -> ConfigMgr: createAdminClient()
activate ConfigMgr
ConfigMgr --> Reader: AdminClient
deactivate ConfigMgr

Reader -> Admin: describeTopics(inputTopics)
activate Admin
Admin -> Kafka: запрос описания топиков
Kafka --> Admin: TopicDescription[]
Admin --> Reader: Map<String, TopicDescription>
deactivate Admin

alt checkACLSWhenTestConnection = true
    Reader -> Reader: checkACLSPermission(desc, READ)
    Reader -> Reader: проверка authorizedOperations
    alt ACL не содержат READ
        Reader --> Framework: ConnectorSecurityException
        deactivate Reader
        deactivate Connector
    end
end

alt executeProbeReadWhenTestConnection = true
    Reader -> Reader: probeRead(topicName)
    activate Reader
    Reader -> ConfigMgr: createConsumer()
    ConfigMgr --> Reader: Consumer
    Reader -> Consumer: assign(partition)
    Reader -> Consumer: poll(Duration)
    Consumer -> Kafka: запрос сообщений
    Kafka --> Consumer: ConsumerRecords
    deactivate Reader
end

Reader -> Admin: describeTopics(outputTopics)
Admin -> Kafka: запрос описания выходных топиков
Kafka --> Admin: TopicDescription[]
Admin --> Reader: Map<String, TopicDescription>

alt checkACLSWhenTestConnection = true
    Reader -> Reader: checkACLSPermission(desc, WRITE)
    alt ACL не содержат WRITE
        Reader --> Framework: ConnectorSecurityException
        deactivate Reader
        deactivate Connector
    end
end

Reader --> Connector: Map<String, TopicDescription>
deactivate Reader
Connector --> Framework: OK (LOGGER.info)
deactivate Connector
@enduml

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

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

participant "IDM Framework" as Framework
participant "KafkaConnector" as Connector
participant "ExecuteQueryHandler" as ExecHandler
participant "TopicReader" as Reader
participant "KafkaConfigManager" as ConfigMgr
participant "Consumer" as Consumer
participant "Kafka Broker" as Kafka
participant "JsonToConnectorObjectMapper" as Mapper
participant "ResultsHandler" as Results
participant "Producer" as Producer

Framework -> Connector: executeQuery(objectClass, filter, resultsHandler, options)
activate Connector
Connector -> ExecHandler: executeQuery(objectClass, filter, resultsHandler, schemaMap)
activate ExecHandler

ExecHandler -> Reader: getConsumerRecords(filter)
activate Reader

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

alt filter != null
    Reader -> Reader: applyFilter(partitions, filter)
    alt AndFilter
        Reader -> Reader: applyAndFilter()
    else GreaterThanOrEqualFilter
        Reader -> Reader: applyGreaterThanOrEqualFilter(PARTITION)
    else LessThanFilter
        Reader -> Reader: applyLessThanFilter(PARTITION)
    else EqualsFilter
        Reader -> Reader: applyEqualFilter(SOURCE)
    end
end

loop для каждой partition
    Reader -> ConfigMgr: createConsumer()
    ConfigMgr --> Reader: Consumer
    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: запрос сообщений
    Kafka --> Consumer: ConsumerRecords
    Reader -> Reader: добавить в результат
    
    alt commitStrategy == COMMIT_AFTER_READ_PARTITION
        Reader -> Consumer: commitSync()
    end
end

Reader --> ExecHandler: List<ConsumerRecord>
deactivate Reader

loop для каждого ConsumerRecord
    ExecHandler -> Mapper: createConnectorObject(record, schemaMap)
    activate Mapper
    Mapper -> Mapper: распарсить JSON из record.value()
    Mapper -> Mapper: преобразовать в ConnectorObject
    Mapper --> ExecHandler: ConnectorObject
    deactivate Mapper
    
    alt успех
        ExecHandler -> Results: handle(connectorObject)
    else исключение
        ExecHandler -> ExecHandler: логировать ошибку
        alt deadLetterQueue настроен
            ExecHandler -> Producer: sendMessageToDeadLetterQueue()
            Producer -> Kafka: отправить в DLQ topic
        end
    end
end

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

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

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

participant "IDM Framework" as Framework
participant "KafkaConnector" as Connector
participant "ScriptHandler" as ScriptH
participant "KafkaExecutorHandler" as ExecutorH
participant "JsonSchema" as JsonSchema
participant "AdminClient" as Admin
participant "Consumer" as Consumer
participant "Producer" as Producer
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, configMgr, reader)
    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 -> Producer: createProducer()
            ExecutorH -> Producer: send(ProducerRecord)
            Producer -> Kafka: отправить сообщение
            Kafka --> Producer: acknowledgment
        else command == "GET_OFFSET"
            ExecutorH -> ExecutorH: getOffset()
            ExecutorH -> Admin: listConsumerGroupOffsets(groupId)
            Admin -> Kafka: запрос offset
            Kafka --> Admin: OffsetAndMetadata
            Admin --> ExecutorH: offset
        else command == "SEEK"
            ExecutorH -> ExecutorH: seek()
            ExecutorH -> Admin: alterConsumerGroupOffsets(groupId, offset)
            Admin -> Kafka: установка offset
            Kafka --> Admin: OK
        else command == "READ_MESSAGES"
            ExecutorH -> ExecutorH: readMessages()
            ExecutorH -> Consumer: getConsumerRecords(filter)
            Consumer -> Kafka: запрос сообщений
            Kafka --> Consumer: ConsumerRecords
            Consumer --> ExecutorH: List<String>
        else
            ExecutorH --> Framework: ConfigurationException
        end
    end
    
    ExecutorH --> ScriptH: result
    ExecutorH -> ExecutorH: dispose()
    ExecutorH -> Admin: close()
    deactivate ExecutorH
else
    ScriptH --> Framework: ConnectorException("Script language not supported")
    deactivate ScriptH
    deactivate Connector
end

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

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

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

participant "ExecuteQueryHandler" as ExecHandler
participant "TopicReader" as Reader
participant "KafkaConfigManager" as ConfigMgr
participant "Consumer" as Consumer
participant "Kafka Broker" as Kafka

ExecHandler -> Reader: getConsumerRecords(filter)
activate Reader

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

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

loop для каждой partition
    Reader -> ConfigMgr: createConsumer()
    ConfigMgr --> Reader: Consumer
    
    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 commitStrategy == COMMIT_AFTER_READ_PARTITION
        Reader -> Consumer: commitSync()
        Consumer -> Kafka: commit offset
        Kafka --> Consumer: OK
    end
end

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

Отправка в Dead Letter Queue#

@startuml
title Диаграмма последовательности: Отправка в Dead Letter Queue

participant "ExecuteQueryHandler" as ExecHandler
participant "KafkaConfigManager" as ConfigMgr
participant "Producer" as Producer
participant "ObjectMapper" as ObjectMapper
participant "Kafka Broker" as Kafka

ExecHandler -> ExecHandler: exception при создании ConnectorObject
ExecHandler -> Config: getDeadLetterQueueOutputTopic()

alt DLQ topic настроен
    ExecHandler -> ConfigMgr: createProducer()
    ConfigMgr --> ExecHandler: Producer
    
    ExecHandler -> ObjectMapper: writeValueAsString(body)
    ExecHandler -> ExecHandler: body = {\n  "badMessageBody": consumerRecordValue,\n  "exceptionMessage": exceptionMessage\n}
    
    ExecHandler -> Producer: send(ProducerRecord)
    Producer -> Kafka: отправить сообщение в DLQ
    Kafka --> Producer: acknowledgment
    
    ExecHandler -> Producer: close()
else DLQ не настроен
    ExecHandler -> ExecHandler: только логирование
end
@enduml