Tous les produits
Search
Centre de documentation

DataHub:Compatibilité Kafka

Dernière mise à jour :Aug 26, 2026

DataHub est entièrement compatible avec le protocole Apache Kafka. Vous pouvez utiliser des clients Kafka natifs pour lire et écrire des données dans DataHub.

Mappage de Kafka vers DataHub

Types de topic

Kafka et DataHub disposent de mécanismes différents pour la mise à l'échelle des topics. Pour garantir la compatibilité avec le comportement de Kafka, définissez le mode de mise à l'échelle sur ONLY_EXTEND lors de la création d'un topic DataHub. Dans ce mode, vous pouvez uniquement ajouter de nouveaux shards à un topic. Ce mode n'autorise ni le fractionnement ni la fusion des shards, et la suppression de shards n'est pas encore prise en charge.

Nommage des topics

Un nom de topic Kafka correspond à un projet et un topic DataHub, séparés par un point (.). Le mappage suit les règles suivantes :

  • La partie située avant le premier . représente le projet DataHub, et la partie suivante représente le topic DataHub. Par exemple, test_project.test_topic correspond au projet test_project et au topic test_topic.

  • Si un nom contient plusieurs caractères ., seul le premier . sert de séparateur. Tous les autres caractères . et - sont remplacés par _.

Partitions

Chaque shard actif dans DataHub correspond à une partition dans Kafka. Par exemple, si un topic comporte cinq shards actifs, il équivaut à un topic Kafka doté de cinq partitions. Lors de l'écriture des données, spécifiez un ID de partition compris dans la plage [0, 4]. Si vous ne spécifiez aucune partition, le client Kafka en attribue une automatiquement.

Topic Tuple

Lorsque vous écrivez des données depuis Kafka vers un topic Tuple, le schéma du topic doit comporter une ou deux colonnes, toutes deux de type STRING. Dans le cas contraire, l'opération d'écriture échoue.

  • Si le schéma comporte une seule colonne, seule la valeur est écrite et la clé est ignorée.

  • Si le schéma comporte deux colonnes, la première et la deuxième colonne correspondent respectivement à la clé et à la valeur.

De plus, n'écrivez pas de données binaires dans un topic Tuple, car cela génère des caractères illisibles. Pour stocker des données binaires, utilisez un topic Blob.

Topic Blob

Lorsque vous écrivez des données depuis Kafka vers un topic Blob, la valeur du message Kafka est écrite dans le champ Blob. Si la clé du message n'est pas NULL, elle est écrite en tant qu'attribut DataHub. Le nom de l'attribut est __kafka_key__ et sa valeur correspond à la clé du message Kafka.

En-têtes

Les en-têtes Kafka correspondent aux attributs DataHub. Si la valeur d'un en-tête est NULL, elle est ignorée et n'est pas écrite en tant qu'attribut. Nous vous recommandons de ne pas utiliser __kafka_key__ comme clé d'en-tête afin d'éviter les conflits avec le nom d'attribut intégré dans les topics Blob.

Groupes de consommateurs

Dans DataHub, un ID d'abonnement fait office de groupe de consommateurs, mais il ne peut s'abonner qu'à un seul topic. En revanche, un groupe de consommateurs Kafka peut s'abonner simultanément à plusieurs topics. Pour assurer la compatibilité avec le modèle d'abonnement de Kafka, DataHub propose une fonctionnalité de groupe. Créez un groupe au sein d'un projet et liez-le à plusieurs topics, ce qui vous permet de vous abonner à tous ces topics via un seul groupe.

Un groupe gère en interne plusieurs abonnements DataHub sur le serveur. Après avoir lié un topic, le groupe crée automatiquement un abonnement, qui apparaît dans la liste des abonnements sur la page de détails du topic. Ne supprimez pas cet abonnement manuellement. Cela empêcherait le groupe de s'abonner au topic et entraînerait la perte de tous les offsets de consommation existants.

Un seul groupe peut s'abonner à un maximum de 50 topics. Pour vous abonner à davantage de topics, soumettez un ticket.

Paramètres Kafka

C=Consommateur, P=Producteur, S=Streams

Paramètre

C/P/S

Valeur

Obligatoire

Description

bootstrap.servers

*

Consultez la section Points de terminaison Kafka.

Oui

security.protocol

*

SASL_SSL

Oui

Pour garantir la sécurité des données, les connexions de Kafka vers DataHub utilisent le chiffrement SSL par défaut.

sasl.mechanism

*

PLAIN

Oui

Le mécanisme d'authentification pour les identifiants AccessKey. Seul PLAIN est pris en charge.

compression.type

P

LZ4

Non

Spécifie le type de compression pour les messages. Actuellement, seul LZ4 est pris en charge.

group.id

C

project.topic:subId

ou

project.group

Oui

Si vous utilisez le format project.topic:subId, l'ID doit correspondre au topic abonné. Sinon, les données ne peuvent pas être lues. Nous vous recommandons d'utiliser le format project.group.

partition.assignment.strategy

C

org.apache.kafka.clients.consumer.RangeAssignor

Non

La stratégie d'attribution des partitions par défaut dans Kafka est RangeAssignor. DataHub ne prend actuellement en charge que cette stratégie. Ne modifiez pas ce paramètre.

session.timeout.ms

C/S

[60000, 180000]

Non

La valeur par défaut dans Kafka est de 10 000 ms. Toutefois, comme DataHub requiert un minimum de 60 000 ms, cette valeur est automatiquement ajustée à 60 000 ms.

heartbeat.interval.ms

C/S

Recommandé : 2/3 de session.timeout.ms

Non

La valeur par défaut de Kafka est de 3 000 ms. Étant donné que session.timeout.ms est ajusté à 60 000 ms, nous vous recommandons de définir explicitement cette valeur sur 40000 afin d'éviter des requêtes de heartbeat trop fréquentes.

application.id

S

project.topic:subId

ou

project.group

Oui

Si vous utilisez le format project.topic:subId, l'ID doit correspondre au topic abonné. Sinon, les données ne peuvent pas être lues. Nous vous recommandons d'utiliser le format project.group.

Ce tableau répertorie les principaux paramètres à examiner lors de l'utilisation d'un client Kafka avec DataHub. Les autres paramètres côté client, tels que retries et batch.size, se comportent comme dans Kafka natif. Les paramètres côté serveur ne modifient pas le comportement réel de DataHub. Par exemple, quelle que soit la valeur de acks, DataHub renvoie une confirmation uniquement après l'écriture complète des données.

Points de terminaison Kafka

Région

ID de région

Point de terminaison public

Point de terminaison ECS (réseau classique)

Point de terminaison ECS (VPC)

Chine (Hangzhou)

cn-hangzhou

dh-cn-hangzhou.aliyuncs.com:9092

dh-cn-hangzhou.aliyun-inc.com:9093

dh-cn-hangzhou-int-vpc.aliyuncs.com:9094

Chine (Shanghai)

cn-shanghai

dh-cn-shanghai.aliyuncs.com:9092

dh-cn-shanghai.aliyun-inc.com:9093

dh-cn-shanghai-int-vpc.aliyuncs.com:9094

Chine (Pékin)

cn-beijing

dh-cn-beijing.aliyuncs.com:9092

dh-cn-beijing.aliyun-inc.com:9093

dh-cn-beijing-int-vpc.aliyuncs.com:9094

Chine (Zhangjiakou)

cn-zhangjiakou

dh-cn-zhangjiakou.aliyuncs.com:9092

dh-cn-zhangjiakou.aliyun-inc.com:9093

dh-cn-zhangjiakou-int-vpc.aliyuncs.com:9094

Chine (Shenzhen)

cn-shenzhen

dh-cn-shenzhen.aliyuncs.com:9092

dh-cn-shenzhen.aliyun-inc.com:9093

dh-cn-shenzhen-int-vpc.aliyuncs.com:9094

Singapour

ap-southeast-1

dh-ap-southeast-1.aliyuncs.com:9092

dh-ap-southeast-1.aliyun-inc.com:9093

dh-ap-southeast-1-int-vpc.aliyuncs.com:9094

Malaisie (Kuala Lumpur)

ap-southeast-3

dh-ap-southeast-3.aliyuncs.com:9092

dh-ap-southeast-3.aliyun-inc.com:9093

dh-ap-southeast-3-int-vpc.aliyuncs.com:9094

Allemagne (Francfort)

eu-central-1

dh-eu-central-1.aliyuncs.com:9092

dh-eu-central-1.aliyun-inc.com:9093

dh-eu-central-1-int-vpc.aliyuncs.com:9094

Chine Est 2 Finance

cn-shanghai-finance-1

dh-cn-shanghai-finance-1.aliyuncs.com:9092

dh-cn-shanghai-finance-1.aliyun-inc.com:9093

dh-cn-shanghai-finance-1-int-vpc.aliyuncs.com:9094

Chine (Hong Kong)

cn-hongkong

dh-cn-hongkong.aliyuncs.com:9092

dh-cn-hongkong.aliyun-inc.com:9093

dh-cn-hongkong-int-vpc.aliyuncs.com:9094

Créer un topic

  1. Créer un topic dans la console

    Lors de la création du topic, activez le mode d'extension des shards.

  2. Créer un topic à l'aide du SDK

    Vous ne pouvez pas créer de topics à l'aide de l'API Kafka. Utilisez le SDK DataHub et définissez ExpandMode sur ONLY_EXTEND. La version de dépendance Maven requise est 2.19.0 ou ultérieure.

    Nous vous recommandons d'utiliser des variables d'environnement pour configurer votre AccessKey ID et votre AccessKey Secret, et d'éviter de les coder en dur dans le code de votre projet. Une paire AccessKey associée à un compte Alibaba Cloud dispose d'autorisations pour toutes les opérations API. Pour une meilleure sécurité, utilisez une paire AccessKey issue d'un utilisateur RAM pour l'accès API ou les opérations quotidiennes afin de réduire le risque de fuite d'identifiants.

    datahub.endpoint=<yourEndpoint>
    datahub.accessId=<yourAccessKeyId>
    datahub.accessKey=<yourAccessKeySecret>
    <dependency>
      <groupId>com.aliyun.datahub</groupId>
      <artifactId>aliyun-sdk-datahub</artifactId>
      <version>2.19.0-public</version>
    </dependency>
    @Value("${datahub.endpoint}")
    String endpoint ;
    @Value("${datahub.accessId}")
    String accessId;
    @Value("${datahub.accessKey}")
    String accessKey;
    public class CreateTopic {
        public static void main(String[] args) {
            DatahubClient datahubClient = DatahubClientBuilder.newBuilder()
                    .setDatahubConfig(
                            new DatahubConfig(endpoint,
                                    new AliyunAccount(accessId, accessKey)))
                    .build();
    
            int shardCount = 1;
            int lifeCycle = 7;
    
            try {
                datahubClient.createTopic("test_project", "test_topic", shardCount, lifeCycle, RecordType.BLOB, "comment", ExpandMode.ONLY_EXTEND);
            } catch (DatahubClientException e) {
                e.printStackTrace();
            }
        }
    }

Créer un groupe

  1. Créer un groupe dans la console

    Cliquez sur Create Group, puis ajoutez les topics auxquels vous souhaitez vous abonner à partir de la liste située à droite. Vous pouvez modifier les topics liés après la création du groupe. Le groupe crée automatiquement un abonnement, qui apparaît sur la page de liste des abonnements du topic.

  2. Créer un groupe à l'aide du SDK

    La version de dépendance Maven doit être 2.21.6-public ou ultérieure.

    <dependency>
      <groupId>com.aliyun.datahub</groupId>
      <artifactId>aliyun-sdk-datahub</artifactId>
      <version>2.21.6-public</version>
    </dependency>
    @Value("${datahub.endpoint}")
    String endpoint ;
    @Value("${datahub.accessId}")
    String accessId;
    @Value("${datahub.accessKey}")
    String accessKey;
    public class CreateGroup {
        public static void main(String[] args) {
            DatahubClient datahubClient = DatahubClientBuilder.newBuilder()
                    .setDatahubConfig(
                            new DatahubConfig(endpoint,
                                    new AliyunAccount(accessId, accessKey)))
                    .build();
    
            List<String> topicList = new ArrayList<>();
            topicList.add("test_project.topic1");
            topicList.add("test_project.topic2");
            topicList.add("test_project.topic3");
    
            try {
                // Create a Kafka group.
                datahubClient.createKafkaGroup("test_project", "test_topic", "test comment");
    
                // Bind the topics to the group for subscription.
                datahubClient.updateTopicsForKafkaGroup("test_project", "test_topic", topicList, UpdateKafkaGroupMode.ADD);
            } catch (DatahubClientException e) {
                e.printStackTrace();
            }
        }
    }

Exemple de producteur

Fichier kafka_client_producer_jaas.conf

Créez un fichier nommé kafka_client_producer_jaas.conf dans n'importe quel répertoire et ajoutez le contenu suivant.

KafkaClient {
  org.apache.kafka.common.security.plain.PlainLoginModule required
  username="yourAccessKeyId"
  password="yourAccessKeySecret";
};

Dépendance Maven

La version du client Kafka doit être 0.10.0.0 ou ultérieure. Nous vous recommandons la version 2.4.0.

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>2.4.0</version>
</dependency>

Exemple de code

public class ProducerExample {
    static {
        System.setProperty("java.security.auth.login.config", "src/main/resources/kafka_client_producer_jaas.conf");
    }

    public static void main(String[] args) {

        Properties properties = new Properties();
        properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
        properties.put("security.protocol", "SASL_SSL");
        properties.put("sasl.mechanism", "PLAIN");
        properties.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        properties.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        properties.put("compression.type", "lz4");

        String KafkaTopicName = "test_project.test_topic";
        Producer<String, String> producer = new KafkaProducer<String, String>(properties);

        try {
            List<Header> headers = new ArrayList<>();
            RecordHeader header1 = new RecordHeader("key1", "value1".getBytes());
            RecordHeader header2 = new RecordHeader("key2", "value2".getBytes());
            headers.add(header1);
            headers.add(header2);

            ProducerRecord<String, String> record = new ProducerRecord<>(KafkaTopicName, 0, "key", "Hello DataHub!", headers);

            // Sync send
            producer.send(record).get();

        } catch (InterruptedException e) {
            e.printStackTrace();
        } catch (ExecutionException e) {
            e.printStackTrace();
        } finally {
            producer.close();
        }
    }
}

Résultat

Après l'exécution réussie du code, prélevez des données pour vérifier le résultat.

Exemple de consommateur

Pour savoir comment générer le fichier kafka_client_producer_jaas.conf et ajouter la dépendance Maven, consultez l'exemple de producteur.

Lorsqu'un nouveau consommateur rejoint le groupe, l'attribution des shards prend entre 10 et 20 secondes. Une fois l'attribution terminée, le consommateur peut commencer à consommer les données.

Exemple de code

Utilisation d'un groupe Kafka (recommandé)

package com.aliyun.datahub.kafka.demo;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;

import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.Properties;
public class ConsumerExample2 {

    static {
        System.setProperty("java.security.auth.login.config", "src/main/resources/kafka_client_producer_jaas.conf");
    }

    public static void main(String[] args) {

        Properties properties = new Properties();
        properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
        properties.put("security.protocol", "SASL_SSL");
        properties.put("sasl.mechanism", "PLAIN");
        // Set group.id to the project.group format.
        properties.put("group.id", "test_project.test_kafka_group");
        properties.put("auto.offset.reset", "earliest");
        properties.put("session.timeout.ms", "60000");
        properties.put("heartbeat.interval.ms", "40000");
        properties.put("ssl.endpoint.identification.algorithm", "");
        properties.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        properties.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        KafkaConsumer<String, String> kafkaConsumer = new KafkaConsumer<String, String>(properties);

        List<String> topicList = new ArrayList<>();
        topicList.add("test_project.test_topic1");
        topicList.add("test_project.test_topic2");
        topicList.add("test_project.test_topic3");
        // By using a Kafka group, you can subscribe to multiple topics.
        kafkaConsumer.subscribe(topicList);

        while (true) {
            ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofSeconds(5));

            for (ConsumerRecord<String, String> record : records) {
                System.out.println(record.toString());
            }
        }
    }
}

Utilisation de project.topic:subId

package com.aliyun.datahub.kafka.demo;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;
public class ConsumerExample {

    static {
        System.setProperty("java.security.auth.login.config", "src/main/resources/kafka_client_producer_jaas.conf");
    }

    public static void main(String[] args) {

        Properties properties = new Properties();
        properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
        properties.put("security.protocol", "SASL_SSL");
        properties.put("sasl.mechanism", "PLAIN");
        // Set group.id to the project.topic:subId format.
        properties.put("group.id", "test_project.test_topic:1611039998153N71KM");
        properties.put("auto.offset.reset", "earliest");
        properties.put("session.timeout.ms", "60000");
        properties.put("heartbeat.interval.ms", "40000");
        properties.put("ssl.endpoint.identification.algorithm", "");
        properties.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        properties.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        KafkaConsumer<String, String> kafkaConsumer = new KafkaConsumer<String, String>(properties);

        // When using the project.topic:subId format, you can only subscribe to a single topic.
        kafkaConsumer.subscribe(Collections.singletonList("test_project.test_topic"));

        while (true) {
            ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofSeconds(5));

            for (ConsumerRecord<String, String> record : records) {
                System.out.println(record.toString());
            }
        }
    }
}

Résultat

Après l'exécution réussie du code, vous pouvez voir les données consommées dans votre terminal.

ConsumerRecord(topic = test_project.test_topic, partition = 0, leaderEpoch = 0, offset = 0, LogAppendTime = 1611040892661, serialized key size = 3, serialized value size = 14, headers = RecordHeaders(headers = [RecordHeader(key = key1, value = [118, 97, 108, 117, 101, 49]), RecordHeader(key = key2, value = [118, 97, 108, 117, 101, 50])], isReadOnly = false), key = key, value = Hello DataHub!)

Dans cet exemple, tous les enregistrements de données renvoyés dans une seule requête partagent le même LogAppendTime, qui correspond à l'horodatage le plus récent parmi tous les enregistrements de ce lot.

Exemple Streams

Dépendance Maven

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>2.4.0</version>
</dependency>
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-streams</artifactId>
    <version>2.4.0</version>
</dependency>

Exemple de code

Cet exemple lit les données d'un topic d'entrée au sein de test_project, convertit les chaînes de clés et de valeurs en minuscules, et écrit le résultat dans un topic de sortie.

public class StreamExample {

    static {
        System.setProperty("java.security.auth.login.config", "src/main/resources/kafka_client_producer_jaas.conf");
    }

    public static void main(final String[] args) {
        final String input = "test_project.input";
        final String output = "test_project.output";
        final Properties properties = new Properties();
        properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
        properties.put("application.id", "test_project.input:1611293595417QH0WL");
        properties.put("security.protocol", "SASL_SSL");
        properties.put("sasl.mechanism", "PLAIN");
        properties.put("session.timeout.ms", "60000");
        properties.put("heartbeat.interval.ms", "40000");
        properties.put("auto.offset.reset", "earliest");

        final StreamsBuilder builder = new StreamsBuilder();
        TestMapper testMapper = new TestMapper();
        builder.stream(input, Consumed.with(Serdes.String(), Serdes.String()))
                .map(testMapper)
                .to(output, Produced.with(Serdes.String(), Serdes.String()));

        final KafkaStreams streams = new KafkaStreams(builder.build(), properties);
        final CountDownLatch latch = new CountDownLatch(1);

        Runtime.getRuntime().addShutdownHook(new Thread("streams-shutdown-hook") {
            @Override
            public void run() {
                streams.close();
                latch.countDown();
            }
        });

        try {
            streams.start();
            latch.await();
        } catch (final Throwable e) {
            System.exit(1);
        }
        System.exit(0);
    }

    static class TestMapper implements KeyValueMapper<String, String, KeyValue<String, String>> {

        @Override
        public KeyValue<String, String> apply(String s, String s2) {
            return new KeyValue<>(StringUtils.lowerCase(s), StringUtils.lowerCase(s2));
        }
    }
}

Résultat

Après le démarrage de la tâche Streams, l'attribution des shards prend environ une minute. Ensuite, vous pouvez voir le nombre de tâches actuelles dans la console. Le nombre de tâches correspond au nombre de shards dans le topic d'entrée. Dans cet exemple, le topic d'entrée comporte trois shards.

currently assigned active tasks: [0_0, 0_1, 0_2]
currently assigned standby tasks: []
revoked active tasks: []
  revoked standby tasks: []

Une fois les shards attribués, écrivez des données de test telles que (AAAA,BBBB),(CCCC,DDDD),(EEEE,FFFF) dans le topic d'entrée. Ensuite, prélevez des données du topic de sortie pour vérifier qu'elles ont été correctement écrites.

Remarques d'utilisation

  • Les transactions et l'idempotence ne sont pas prises en charge.

  • Les clients Kafka ne peuvent pas créer automatiquement des topics dans DataHub. Vous devez créer le topic avant d'y écrire des données.

  • Lors de l'utilisation d'un ID d'abonnement (project.topic:subid) comme group.id, un consommateur ne peut s'abonner qu'à un seul topic. Pour vous abonner à plusieurs topics, utilisez un groupe DataHub.

  • L'horodatage des données lues par un consommateur correspond toujours au LogAppendTime, qui indique le moment où les données ont été écrites dans DataHub. Tous les enregistrements d'une seule requête de récupération partagent le même horodatage : l'horodatage le plus récent de ce lot. Cela signifie que l'horodatage de lecture peut être postérieur à l'heure d'écriture réelle.

  • Une application Streams ne prend en charge qu'un seul topic d'entrée, mais peut avoir plusieurs topics de sortie.

  • Seules les tâches Streams sans état sont prises en charge.

  • Les versions Kafka prises en charge vont de 0.10.0 à 2.4.0.

FAQ

Déconnexion lors de l'écriture des données

Selector - [Producer clientId=producer-1] Connection with dh-cn-shenzhen.aliyuncs.com disconnected
java.io.EOFException
    at org.apache.kafka.common.network.SslTransportLayer.read(SslTransportLayer.java:573)
    ...

Les requêtes de métadonnées Kafka et les requêtes d'écriture de données utilisent des connexions différentes.

Le client établit d'abord une connexion pour récupérer les métadonnées. Il utilise ensuite les informations de broker renvoyées pour établir une seconde connexion destinée à l'écriture des données. Toutes les requêtes suivantes sont envoyées via cette seconde connexion.

La première connexion, désormais inactive, est automatiquement fermée par le serveur après un délai d'expiration. Cela peut générer une erreur de déconnexion dans les journaux. Vous pouvez ignorer cette erreur si les données sont écrites avec succès.

Échec du démarrage du client Kafka

Caused by: org.apache.kafka.common.errors.SslAuthenticationException: SSL handshake failed
Caused by: javax.net.ssl.SSLHandshakeException: No subject alternative names matching IP address 100.67.134.161 found

Ajoutez la propriété suivante à votre configuration : properties.put("ssl.endpoint.identification.algorithm", "");.

DisconnectException lors de la consommation

[INFO][Consumer clientId=client-id, groupId=consumer-project.topic:subid] Error sending fetch request (sessionId=INVALID, epoch=INITIAL) to node 1: {}.
org.apache.kafka.common.errors.DisconnectException

Le client Kafka doit maintenir une connexion TCP persistante avec le serveur. Cette exception est généralement causée par des fluctuations réseau. Le client dispose d'une logique de nouvelle tentative intégrée, de sorte que cette erreur n'affecte généralement pas la consommation.