Pratique de l'utilisation hors ligne de Flink basé sur TDMQ Apache Pulsar

Apache Flink est un framework open source pour le traitement en continu et par lots, doté d'un moteur de streaming à haut débit et faible latence, prenant en charge le traitement par temps d'événement, la gestion d'état et assurant la tolérance aux pannes ainsi que la sémantique exactly-once. Le cœur de Flink est un moteur distribué de traitement de flux de données, supportant les langages Java, Scala, Python et SQL, pouvant exécuter des programmes de flux de données dans des environnements cluster ou cloud. Il propose l'API DataStream pour traiter des flux de données bornés ou non bornés, l'API DataSet pour les jeux de données bornés, ainsi que les interfaces Table API et SQL pour le traitement relationnel en continu et par lots. Actuellement, Flink a été mis à jour jusqu'à la version 1.20, et durant cette évolution, non seulement le framework lui-même mais aussi les plugins ont subi des changements dans leurs APIs et configurations. Cet article se concentre principalement sur les tests et validations du plugin Pulsar Flink pour les versions récentes, notamment la version 1.17. La version actuelle de Flink peut être consultée ici : https://nightlies.apache.org/flink/

Déploiement de Flink

Configuration de l'environnnement Flink

Suivant la documentation officielle de Flink 1.17, le déploiement de la version Docker de Flink peut être effectué via ce lien : https://nightlies.apache.org/flink/flink-docs-release-1.17/docs/deployment/resource-providers/standalone/docker/#getting-started
Tout d'abord, configurez les informations d'environnement du JobManager et TaskManager du cluster Flink. Notez que le connecteur Pulsar utilise de la mémoire hors tas, et que la mémoire hors tas par défaut pour les tâches est définie à zéro. Par conséquent, il est nécessaire de spécifier explicitement la taille de la mémoire hors tas via taskmanager.memory.task.off-heap.size, ici fixée à 1 Go : https://nightlies.apache.org/flink/flink-docs-release-1.17/docs/deployment/memory/mem_setup_tm/#configure-off-heap-memory-direct-or-native

$ FLINK_PROPERTIES=$'\njobmanager.rpc.address: jobmanager\ntaskmanager.memory.task.offheap.size: 1gb\ntaskmanager.memory.process.size: 4gb'
$ docker network create flink-network

Déploiement du JobManager

Après configuraton des variables d'environnement, déployez le JobManager. Le port par défaut est 8081. Une fois déployé, vous pouvez accéder au tableau de bord Flink via le port 8081.

$ docker run \
    --rm \
    --name=jobmanager \
    --network flink-network \
    --publish 8081:8081 \
    --env FLINK_PROPERTIES="${FLINK_PROPERTIES}" \
    flink:1.17.2-scala_2.12 jobmanager

Déploiement du TaskManager

Le JobManager est responsable de la coordination des tâches. Après avoir déployé le JobManager, il faut également déployer le TaskManager qui exécute les tâches.

$ docker run \
    --rm \
    --name=taskmanager \
    --network flink-network \
    --env FLINK_PROPERTIES="${FLINK_PROPERTIES}" \
    flink:1.17.2-scala_2.12 taskmanager

Une fois le TaskManager lancé, vous pouvez voir qu'il est correctement enregistré dans la console du JobManager 8081. Le déploiement des composants Docker Flink est maintenant terminé.

Téléchargement de Flink Cli

Après avoir compilé et empaqueté votre tâche Pulsar localement, vous devez utiliser Flink Cli pour soumettre la tâche au cluster Docker Flink. Téléchargez le fichier binaire Flink correspondant à la version Docker depuis ce lien et extrayez-le localement : https://flink.apache.org/downloads/

Démo : Copie de Topics

Suivant la documentation communautaire du connecteur Flink Pulsar et les documents liés à Oceanus, ce démo utilise le SDK Flink 1.17 pour copier tous les messages d'un topic dans un namespace vers un autre topic. Ce démo illustre l'utilisation basique du connecteur Flink sans utiliser de sérialiseurs personnalisés, mais plutôt en utilisant le sérialiseur String intégré au connecteur. https://cloud.tencent.com/document/product/849/85885#pulsar-source-.E5.92.8C-sink-.E7.A4.BA.E4.BE.8B https://nightlies.apache.org/flink/flink-docs-release-1.17/docs/connectors/datastream/pulsar/#apache-pulsar-connector

Logique principale

La logique principale est illustrée dans le code ci-dessous. Utilisez d'abord l'outil ParameterTool pour analyser les paramètres passés en ligne de commande. Ensuite, créez un flux Flink Stream à partir du topic d'entrée vers le topic de sortie en utilisant les méthodes Builder du Source et Sink du connecteur.

public static void main(String[] args) throws Exception {
    final ParameterTool parameterTool = ParameterTool.fromArgs(args);
    if (parameterTool.getNumberOfParameters() < 2) {
        System.err.println("Missing parameters!");
        return;
    }
    final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    env.getConfig().setRestartStrategy(RestartStrategies.fixedDelayRestart(4, 10000));
    env.enableCheckpointing(60000);
    env.getConfig().setGlobalJobParameters(parameterTool);
    String brokerServiceUrl = parameterTool.getRequired("broker-service-url");
    String inputTopic = parameterTool.getRequired("input-topic");
    String outputTopic = parameterTool.getRequired("output-topic");
    String subscriptionName = parameterTool.get("subscription-name", "testDuplicate");
    String token = parameterTool.getRequired("token");
    // source
    PulsarSource<String> source = PulsarSource.builder()
            .setServiceUrl(brokerServiceUrl)
            .setStartCursor(StartCursor.latest())
            .setTopics(inputTopic)
            .setDeserializationSchema(new SimpleStringSchema())
            .setSubscriptionName(subscriptionName)
            .setAuthentication("org.apache.pulsar.client.impl.auth.AuthenticationToken", token)
            .build();
    DataStream<String> stream = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Pulsar Source");
    // sink
    PulsarSink<String> sink = PulsarSink.builder()
            .setServiceUrl(brokerServiceUrl)
            .setTopics(outputTopic)
            .setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE)
            .setAuthentication("org.apache.pulsar.client.impl.auth.AuthenticationToken", token) 
            .setConfig(PulsarSinkOptions.PULSAR_BATCHING_ENABLED, false)
            .setSerializationSchema(new SimpleStringSchema())
            .build();
    stream.sinkTo(sink);
    env.execute("Pulsar Streaming Message Duplication");
}

Vérification

Créez les topics d'entrée NinjaDuplicationInput1 et de sortie NinjaDuplicationOutput1 via la console TDMQ Pulsar.

Après compilation du code en jar, téléchargez le jar sur le cluster Docker Flink. Obtenez un token avec les rôles de production et de consommation via l'interface de gestion des rôles. La commande est la suivante :

/usr/local/services/flink/flink-1.17.2 # /usr/local/services/flink/flink-1.17.2/bin/flink run /tmp/wordCount/pulsar-flink-examples-0.0.1-SNAPSHOT-jar-with-dependencies.jar \
    --broker-service-url http://pulsar-xxxxx \
    --input-topic pulsar-g8akj4eow8z8/dev-tdmq-ninjazhou-1713856927/ninjaDuplicationInput1 \
    --outputtopic pulsar-g8akj4eow8z8/dev-tdmq-ninjazhou-1713856927/ninjaDuplicationOutput1 \
    --subscription-name ninjaTest1 \
    --token eyJrZXlJZCI6InB1bHNhci1nOGFrajRlb3c4ejgiLCJhbGciOiJIUzI1NiJ9.eyJzdWIiOiJwdWxzYXItZzhha2o0ZW93OHo4X2Rldi10ZG1xLW5pbmphemhvdS0xNzEzODU2OTI3LWRldmNsb3VkIn0.O89TuTl0PMca7tLN9aGvHvqt7QZ9yMrJh1z3VOz7EVc
Job has been submitted with JobID c1bdab89c01ef16e00579bd2c6648859

Après soumission de la tâche, vous pouvez voir apparaître la tâche dans le tableau de bord Flink avec l'état "Running".

Envoyez des messages sur le topic NinjaDuplicationInput1 via la ligne de commande :

/usr/local/services/pulsar/apache-pulsar-2.9.5/bin/pulsar-client \
--url http://pulsar-xxxxxx \
--auth-params token:eyJrZXlJZCI6InB1bHNhci1nOGFrajRlb3c4ejgiLCJhbGciOiJIUzI1NiJ9.eyJzdWIiOiJwdWxzYXItZzhha2o0ZW93OHo4X2Rldi10ZG1xLW5pbmphemhvdS0xNzEzODU2OTI3LWRldmNsb3VkIn0.O89TuTl0PMca7tLN9aGvHvqt7QZ9yMrJh1z3VOz7EVc \
--auth-plugin org.apache.pulsar.client.impl.auth.AuthenticationToken \
produce \
-m "i am the bone of my sword" \
-n 5 \
pulsar-g8akj4eow8z8/dev-tdmq-ninjazhou-1713856927/ninjaDuplicationInput1

Après l'envoi des messages, vous pouvez observer dans la console de requête de messages que le topic de destination NinjaDuplicationOutput1 contient désormais cinq messages identiques à ceux envoyés.

Les logs de sortie standard du TaskManager Docker montrent également les messages envoyés par le Sink vers le topic cible.

Démo : Comptage des mots

Le comptage des mots est un exemple courant dans Flink, illustrant bien la philosophie de traitement en flux. Ce démo suit l'exemple de StreamNative, utilisant le SDK Flink 1.17 pour traiter un topic Pulsar comme ressource source et destination, comptant le nombre d'apparitions de chaque mot dans une fenêtre temporelle donnée, et envoyant les résultats vers un topic cible. https://github.com/streamnative/examples/blob/master/pulsar-flink/README.md

Logique principale

Le projet complet est disponible ici : pulsar-flink-example.zip. La logique centrale est illustrée dans le code ci-dessous. Utilisez d'abord l'outil ParameterTool pour analyser les paramètres passés en ligne de commande, puis utilisez le désérialiseur intégré de Flink pour convertir le corps du message en chaîne de caractères. Dans la partie de traitement des données, utilisez des fenêtres de temps système pour compter les messages entrants dans une période donnée, et générer des objets WordCount pour chaque mot. Enfin, utilisez un sérialiseur personnalisé pour convertir les objets WordCount en tableau d'octets JSON, envoyés vers le topic cible. Actuellement, le connecteur TDMQ Pulsar prend en charge trois méthodes pour sérialiser des objets Java en messages d'octets pour le Sink Pulsar : Schéma Pulsar, Schéma Flink et sérialiseur personnalisé. Il est recommandé d'utiliser le sérialiseur personnalisé pour les objets WordCount définis. https://nightlies.apache.org/flink/flink-docs-release-1.17/docs/connectors/datastream/pulsar/#serializer

Notez que le mode d'envoi par lot est activé par défaut dans le Sink. Lors de la recherche de messages dans la console, seul le premier message d’un batch est visible, ce qui complique la comparaison du nombre de messages. Le démo désactive cette fonctionnalité.

/**
 * Référence au démo pulsar flink de streamNative
 * <a href="https://github.com/streamnative/examples/tree/master/pulsar-flink">pulsar-flink example</a>
 * Étant donné que le démo streamNative utilise Flink 1.10.1 et le connecteur Pulsar 2.4.17,
 * et que les APIs de Flink et du connecteur Pulsar ont évolué dans la version 1.20 de la communauté,
 * ce démo est réécrit en utilisant Flink 1.17.
 * <a href="https://nightlies.apache.org/flink/flink-docs-release-1.17/docs/deployment/resource-providers/standalone/overview/">Documentation Flink 1.17</a>
 * <p>
 * Ce démo compte la fréquence d’apparition de chaque mot dans une fenêtre temporelle donnée dans le topic source,
 * et envoie les résultats sous forme de messages uniques par mot dans le topic cible.
 *
 */
public class PulsarStreamingWordCount {
    private static final Logger LOG = LoggerFactory.getLogger(PulsarStreamingWordCount.class);
    
    public static void main(String[] args) throws Exception {
        // Analyse des paramètres de tâche
        // Authentification par default avec authToken
        final ParameterTool parameterTool = ParameterTool.fromArgs(args);
        if (parameterTool.getNumberOfParameters() < 2) {
            System.err.println("Missing parameters!");
            return;
        }
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.getConfig().setRestartStrategy(RestartStrategies.fixedDelayRestart(4, 10000));
        env.enableCheckpointing(60000);
        env.getConfig().setGlobalJobParameters(parameterTool);
        String brokerServiceUrl = parameterTool.getRequired("broker-service-url");
        String inputTopic = parameterTool.getRequired("input-topic");
        String outputTopic = parameterTool.getRequired("output-topic");
        String subscriptionName = parameterTool.get("subscription-name", "WordCountTest");
        String token = parameterTool.getRequired("token");
        int timeWindowSecond = parameterTool.getInt("time-window", 60);
        // source
        PulsarSource<String> source = PulsarSource.builder()
                .setServiceUrl(brokerServiceUrl)
                .setStartCursor(StartCursor.latest())
                .setTopics(inputTopic)
                // Ici, le payload du message est sérialisé en chaîne de caractères
                // Le source ne supporte actuellement que le parsing du contenu du payload, pas les propriétés du message
                // Comme publish_time
                // Pour parser les propriétés, il faut implémenter la méthode getProducedType() dans une classe héritée
                // La méthode getProducedType est complexe, nécessitant la déclaration de chaque attribut désérialisé
                // https://nightlies.apache.org/flink/flink-docs-release-1.17/docs/connectors/datastream/pulsar/#deserializer
                .setDeserializationSchema(new SimpleStringSchema())
                .setSubscriptionName(subscriptionName)
                .setAuthentication("org.apache.pulsar.client.impl.auth.AuthenticationToken", token)
                .build();
        // Comme nous n'utilisons pas le temps de publication du message
        // Nous utilisons le mode noWatermark avec le temps système du taskManager
        DataStream<String> stream = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Pulsar Source");
        // process
        // Découpe chaque ligne en mots, puis agrège les occurrences dans un objet WordCount
        // Utilisation de TumblingProcessingTimeWindows pour les fenêtres basées sur le temps système
        DataStream<WordCount> wc = stream
                .flatMap((FlatMapFunction<String, WordCount>) (line, collector) -> {
                    LOG.info("current line = {}, word list = {}", line, line.split("\\s"));
                    for (String word : line.split("\\s")) {
                        collector.collect(new WordCount(word, 1, null));
                    }
                })
                .returns(WordCount.class)
                .keyBy(WordCount::getWord)
                .window(TumblingProcessingTimeWindows.of(Time.seconds(timeWindowSecond)))
                .reduce((ReduceFunction<WordCount>) (c1, c2) -> {
                    WordCount reducedWordCount = new WordCount(c1.getWord(), c1.getCount() + c2.getCount(), null);
                    LOG.info("previous [{}] [{}], current wordCount {}", c1, c2, reducedWordCount);
                    return reducedWordCount;
                });
        // sink
        // La version 1.17 de Flink fournit deux méthodes de sérialisation déjà implémentées :
        // une utilisant le schéma intégré de Pulsar, l'autre le schéma de Flink
        // Mais la version 2.9 de Pulsar de TDMQ ne supporte pas encore parfaitement le schéma
        // On utilise donc l'interface PulsarSerializationSchema<T> de Flink, implémentant la méthode serialize(IN element, PulsarSinkContext sinkContext)
        // pour convertir l'objet IN en tableau d'octets
        // https://nightlies.apache.org/flink/flink-docs-release-1.17/docs/connectors/datastream/pulsar/#serializer
        PulsarSink<WordCount> sink = PulsarSink.builder()
                .setServiceUrl(brokerServiceUrl)
                .setTopics(outputTopic)
                .setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE)
                .setAuthentication("org.apache.pulsar.client.impl.auth.AuthenticationToken", token)
                .setConfig(PulsarSinkOptions.PULSAR_BATCHING_ENABLED, false)
                .setSerializationSchema(new PulsarSerializationSchema<WordCount>() {
                    private ObjectMapper objectMapper;
                    @Override
                    public void open(
                            SerializationSchema.InitializationContext initializationContext,
                            PulsarSinkContext sinkContext,
                            SinkConfiguration sinkConfiguration)
                            throws Exception {
                        objectMapper = new ObjectMapper();
                    }
                    @Override
                    public PulsarMessage<?> serialize(WordCount wordCount, PulsarSinkContext sinkContext) {
                        // Ajout du timestamp de traitement et sérialisation en JSON
                        byte[] wordCountBytes;
                        wordCount.setSinkDateTime(LocalDateTime.now().toString());
                        try {
                            wordCountBytes = objectMapper.writeValueAsBytes(wordCount);
                        } catch (Exception exception) {
                            wordCountBytes = exception.getMessage().getBytes();
                        }
                        return PulsarMessage.builder(wordCountBytes).build();
                    }
                })
                .build();
        wc.sinkTo(sink);
        env.execute("Pulsar Streaming WordCount");
    }
}

Vérification

Créez les topics d'entrée NinjaWordCountInput1 et de sortie NinjaWordCountOutput1 via la console TDMQ Pulsar.

Après compilation du code en jar, téléchargez le jar sur le cluster Docker Flink. Obtenez un token avec les rôles de production et de consommation via l'interface de gestion des rôles. La commande est la suivante :

/usr/local/services/flink/flink-1.17.2 # /usr/local/services/flink/flink-1.17.2/bin/flink run /tmp/wordCount/pulsar-flink-examples-0.0.1-SNAPSHOT-jar-with-dependencies.jar \
--broker-service-url http://pulsar-xxxx \
--input-topic pulsar-g8akj4eow8z8/dev-tdmq-ninjazhou-1713856927/ninjaWordCountInput1 \
--output-topic pulsar-g8akj4eow8z8/dev-tdmq-ninjazhou-1713856927/ninjaWordCountOutput1 \
--subscription-name ninjaTest3 \
--token eyJrZXlJZCI6InB1bHNhci1nOGFrajRlb3c4ejgiLCJhbGciOiJIUzI1NiJ9.eyJzdWIiOiJwdWxzYXItZzhha2o0ZW93OHo4X2Rldi10ZG1xLW5pbmphemhvdS0xNzEzODU2OTI3LWRldmNsb3VkIn0.O89TuTl0PMca7tLN9aGvHvqt7QZ9yMrJh1z3VOz7EVc
Job has been submitted with JobID 6f608d95506f96c3eac012386f840655

Après soumission de la tâche, vous pouvez voir apparaître la tâche dans le tableau de bord Flink avec l'état "Running".

Envoyez des messages sur le topic NinjaWordCountInput1 via la ligne de commande. Deux lots sont envoyés : le premier avec "i am the bone of my sword" 5 fois, le second avec "Test1" 3 fois.

/usr/local/services/pulsar/apache-pulsar-2.9.5/bin/pulsar-client \
--url http://pulsar-xxx \
--auth-params token:eyJrZXlJZCI6InB1bHNhci1nOGFrajRlb3c4ejgiLCJhbGciOiJIUzI1NiJ9.eyJzdWIiOiJwdWxzYXItZzhha2o0ZW93OHo4X2Rldi10ZG1xLW5pbmphemhvdS0xNzEzODU2OTI3LWRldmNsb3VkIn0.O89TuTl0PMca7tLN9aGvHvqt7QZ9yMrJh1z3VOz7EVc \
--auth-plugin org.apache.pulsar.client.impl.auth.AuthenticationToken \
produce \
-m "i am the bone of my sword" \
-n 5 \
pulsar-g8akj4eow8z8/dev-tdmq-ninjazhou-1713856927/ninjaWordCountInput1
/usr/local/services/pulsar/apache-pulsar-2.9.5/bin/pulsar-client \
--url http://pulsar-g8akj4eow8z8.sap-8ywks40k.tdmq.ap-gz.qcloud.tencenttdmq.com:8080 \
--auth-params token:eyJrZXlJZCI6InB1bHNhci1nOGFrajRlb3c4ejgiLCJhbGciOiJIUzI1NiJ9.eyJzdWIiOiJwdWxzYXItZzhha2o0ZW93OHo4X2Rldi10ZG1xLW5pbmphemhvdS0xNzEzODU2OTI3LWRldmNsb3VkIn0.O89TuTl0PMca7tLN9aGvHvqt7QZ9yMrJh1z3VOz7EVc \
--auth-plugin org.apache.pulsar.client.impl.auth.AuthenticationToken \
produce \
-m "test1" \
-n 3 \
pulsar-g8akj4eow8z8/dev-tdmq-ninjazhou-1713856927/ninjaWordCountInput1

Après l'envoi des messages, vous pouvez observer dans la console de requête de messages que le topic de destination NinjaWordCountOutput1 contient désormais 8 messages.

Chaque message est un tableau d'octets JSON contenant le nom du mot, sa fréquence et le point de temps de traitement. Le diagramme montre la structure du message pour le mot "am", confirmant que le nombre d'occurrences correspond au nombre de messages envoyés, prouvant que la tâche fonctionne correctement.

En consultant le TaskManager, vous pouvez voir le contenu des messages ainsi que le processus de parsing de chaque message.

Résumé de l'utilisation du connecteur Flink

Choix de version

Actuellement, pour la production et la consommation avec le connecteur Flink, après investigation, il est possible de répondre aux besoins de base de TDMQ Pulsar sans modification de contrôle ni opération non standard. À ce jour, Apache Flink a publié la version 1.20, et il est recommandé d'utiliser Flink 1.15 à 1.17 avec le connecteur Pulsar. Les versions 1.15 et inférieures ne sont pas recommandées, tandis que les versions 1.18 et supérieures peuvent être utilisées comme la version 1.17.

Voici les principales configurations du connecteur Flink Pulsar pour les versions 1.15 et 1.17. Les dépendances du connecteur Flink correspondantes à chaque version peuvent être trouvées ici : https://nightlies.apache.org/flink/flink-docs-release-1.17/docs/connectors/datastream/pulsar/#dependency

Liens vers les documentations de chaque version : https://nightlies.apache.org/flink/

Connecteur Pulsar Flink 1.17

Dépendances du code

Dans un projet Java, ajoutez les dépendances suivantes dans pom.xml :

<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-base</artifactId>
<version>1.17.2</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-pulsar</artifactId>
<version>4.1.0-1.17</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>1.17.2</version>
</dependency>

Exemple de code Source

PulsarSource<String> source = PulsarSource.builder()
.setServiceUrl(brokerServiceUrl)
.setStartCursor(StartCursor.latest())
.setTopics(inputTopic)
.setDeserializationSchema(new SimpleStringSchema())
.setSubscriptionName(subscriptionName)
.setAuthentication("org.apache.pulsar.client.impl.auth.AuthenticationToken", token)
.build();

Description des paramètres Source

Tous les paramètres du Source du connecteur peuvent être consultés dans la documentation officielle. Voici les paramètres couramment utilisés :

Nom du paramètre Description
setServiceUrl Adresse d'accès TDMQ Pulsar, par exemple http://pulsar-xxx:8080
setStartCursor Point de départ du topic, supporte earliest, latest, ID de message et point de temps. Si un abonnement existe déjà, le point de l'abonnement est utilisé
setTopics Nom du topic, par exemple pulsar-xxxxx/dev-tdmq-ninjazhou-1713856927/ninjaWordCountInput1
setDeserializationSchema Schéma de désérialisation, il est recommandé d'utiliser le sérialiseur SimpleStringSchema intégré à Flink ou StringSchema de Pulsar
setSubscriptionName Nom de l'abonnement
setAuthentication Classe d'authentification. Pour TDMQ Pulsar, utiliser toujours AuthenticationToken avec un jeton JWT. Le jeton doit être un secret de rôle avec droits de lecture sur le topic

Exemple de code Sink

PulsarSink<String> sink = PulsarSink.builder()
.setServiceUrl(brokerServiceUrl)
.setTopics(outputTopic)
.setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE)
.setAuthentication("org.apache.pulsar.client.impl.auth.AuthenticationToken", token)
.setSerializationSchema(new SimpleStringSchema())
.build();

Description des paramètres Sink

Tous les paramètres du Sink du connecteur peuvent être consultés dans la documentation officielle. Voici les paramètres couramment utilisés :

Nom du paramètre Description
setServiceUrl Adresse d'accès TDMQ Pulsar, par exemple http://pulsar-xxx:8080
setTopics Nom du topic, par exemple pulsar-xxxx/dev-tdmq-ninjazhou-1713856927/ninjaWordCountOutput1
setSerializationSchema Sérialiseur transformant l'objet en tableau d'octets. Il est recommandé d'implémenter l'interface PulsarSerializationSchema personnalisée
setDeliveryGuarantee Garantie de livraison, options disponibles : NONE, AT_LEAST_ONCE, EXACTLY_ONCE. Seul AT_LEAST_ONCE ou NONE est recommandé
setAuthentication Classe d'authentification. Pour TDMQ Pulsar, utiliser toujours AuthenticationToken avec un jeton JWT. Le jeton doit être un secret de rôle avec droits d'écriture sur le topic

Connecteur Pulsar Flink 1.15

Dépendances du code

Dans un projet Maven, ajoutez les dépendances suivantes dans pom.xml :

<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-base</artifactId>
<version>1.15.4</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-pulsar</artifactId>
<version>1.15.4</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>1.15.4</version>
</dependency>

Exemple de code Source

PulsarSource<String> source = PulsarSource.builder()
.setServiceUrl(brokerServiceUrl)
.setAdminUrl(brokerServiceUrl)
.setStartCursor(StartCursor.latest())
.setTopics(inputTopic)
.setDeserializationSchema(PulsarDeserializationSchema.flinkSchema(new SimpleStringSchema()))
.setSubscriptionName(subscriptionName)
.setSubscriptionType(SubscriptionType.Exclusive)
.setConfig(PulsarOptions.PULSAR_AUTH_PLUGIN_CLASS_NAME, "org.apache.pulsar.client.impl.auth.AuthenticationToken")
.setConfig(PulsarOptions.PULSAR_AUTH_PARAMS, token)
.setConfig(PulsarSourceOptions.PULSAR_ENABLE_AUTO_ACKNOWLEDGE_MESSAGE, true)
.build();

Description des paramètres Source

Tous les paramètres du Source du connecteur peuvent être consultés dans la documentation officielle. Voici les paramètres couramment utilisés :

Nom du paramètre Description
setServiceUrl Adresse d'accès TDMQ Pulsar, par exemple http://pulsar-xxxxx:8080
setStartCursor Point de départ du topic, supporte earliest, latest, ID de message et point de temps. Si un abonnement existe déjà, le point de l'abonnement est utilisé
setTopics Nom du topic, par exemple pulsar-xxxx/ninjaWordCountInput1
setDeserializationSchema Schéma de désérialisation, il est recommandé d'utiliser le sérialiseur SimpleStringSchema intégré à Flink ou StringSchema de Pulsar
setSubscriptionName Nom de l'abonnement
setConfig(PulsarOptions.PULSAR_AUTH_PLUGIN_CLASS_NAME) Classe d'authentification. Pour TDMQ Pulsar, utiliser toujours AuthenticationToken
setConfig(PulsarOptions.PULSAR_AUTH_PARAMS) Valeur d'authentification. Pour TDMQ Pulsar, utiliser toujours un jeton JWT. Le jeton doit être un secret de rôle avec droits de lecture sur le topic
setAdminUrl Adresse du point d'accès de gestion, requis pour les anciennes versions pour les opérations de gestion
setSubscriptionType Type d'abonnement. Les versions anciennes nécessitent de spécifier Exclusive ou Failover. Shared est déprécié
setConfig(PulsarSourceOptions.PULSAR_ENABLE_AUTO_ACKNOWLEDGE_MESSAGE) Si désactivé, le plugin utilise les transactions pour les accusés de réception. Si activé, les messages sont automatiquement reconnus selon un intervalle

Exemple de code Sink

PulsarSink<String> sink = PulsarSink.builder()
.setServiceUrl(brokerServiceUrl)
.setAdminUrl(brokerServiceUrl)
.setTopics(outputTopic)
.setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE)
.setConfig(PulsarOptions.PULSAR_AUTH_PLUGIN_CLASS_NAME, "org.apache.pulsar.client.impl.auth.AuthenticationToken")
.setConfig(PulsarOptions.PULSAR_AUTH_PARAMS, token)
.setSerializationSchema(PulsarSerializationSchema.flinkSchema(new SimpleStringSchema()))
.build();

Description des paramètres Sink

Tous les paramètres du Sink du connecteur peuvent être consultés dans la documentation officielle. Voici les paramètres couramment utilisés :

Nom du paramètre Description
setServiceUrl Adresse d'accès TDMQ Pulsar, par exemple http://pulsar-xxxx:8080
setTopics Nom du topic, par exemple pulsar-xxxx/dev-tdmq-ninjazhou-1713856927/ninjaWordCountOutput1
setSerializationSchema Sérialiseur transformant l'objet en tableau d'octets. Il est recommandé d'implémenter l'interface PulsarSerializationSchema personnalisée
setDeliveryGuarantee Garantie de livraison, options disponibles : NONE, AT_LEAST_ONCE, EXACTLY_ONCE. Seul AT_LEAST_ONCE ou NONE est recommandé
setAdminUrl Adresse du point d'accès de gestion, requis pour les anciennes versions pour les opérations de gestion
setConfig(PulsarOptions.PULSAR_AUTH_PLUGIN_CLASS_NAME) Classe d'authentification. Pour TDMQ Pulsar, utiliser toujours AuthenticationToken
setConfig(PulsarOptions.PULSAR_AUTH_PARAMS) Valeur d'authentification. Pour TDMQ Pulsar, utiliser toujours un jeton JWT. Le jeton doit être un secret de rôle avec droits d'écriture sur le topic

Points importants à retenir

  1. Le connecteur Pulsar utilise de la mémoire hors tas. La taille de la mémoire hors tas doit être explicitement définie via taskmanager.memory.task.off-heap.size, par exemple 1Go.
  2. Le setSerializationSchema offre deux méthodes implémentées : schéma Pulsar ou schéma Flink. Ces deux méthodes couplent le code métier au schéma. Il est recommandé d'implémenter l'interface PulsarSerializationSchema<IN> et de définir la méthode serialize(IN element, PulsarSinkContext sinkContext).
  3. Le Sink est activé par défaut en mode batch. Pour désactiver ce mode, configurez setConfig(PulsarSinkOptions.PULSAR_BATCHING_ENABLED, false).
  4. Le traitement en fenêtre dans Flink supporte deux types de temps : ProcessTime et EventTime. Cependant, le Source ne peut actuellement parser que le payload du message, pas les propriétés comme publish_time. Pour parser les propriétés, il faut implémenter getProducedType() dans une classe héritée. Il est conseillé d'utiliser ProcessTime.
  5. Dans les versions Flink 1.16 et antérieures, setSubscriptionType supportait Shared et Key_shared. Ces modes nécessitent des transactions, mais TDMQ Pulsar ne supporte pas encore les transactions. Il est donc impossible de combiner setSubscriptionType(SubscriptionType.Shared) et setConfig(PulsarSourceOptions.PULSAR_ENABLE_AUTO_ACKNOWLEDGE_MESSAGE, False).
  6. Le connecteur Pulsar intégré dans Oceanus est basé sur StreamNative et compatible avec Flink 1.13-1.14. Ces versions sont obsolètes et présentent des incompatibilités API avec les nouvelles versions. L'utilisation de ces versions avec Flink récent nécessiterait des adaptations importantes du code.

Étiquettes: Apache Flink TDMQ Apache Pulsar streaming batch processing

Publié le 20 septembre à 07h18