Cette rubrique décrit les principaux paramètres du développement SQL et fournit des explications ainsi que des exemples d'utilisation.
table.exec.sink.keyed-shuffle
Pour résoudre les problèmes de désordre lors de l'écriture de données dans une table dotée d'une clé primaire, utilisez le paramètre table.exec.sink.keyed-shuffle pour effectuer un brassage par hachage (hash shuffle). Cette opération garantit que les enregistrements partageant la même clé primaire sont acheminés vers la même instance d'opérateur, ce qui atténue les problèmes de désordre.
Remarques d'utilisation
Le brassage par hachage n'est efficace que si l'opérateur en amont peut garantir l'ordre des enregistrements de mise à jour pour une même clé primaire. Dans le cas contraire, cette opération ne résout pas le problème de désordre.
Si vous modifiez le parallélisme d'un opérateur en mode expert, les règles de parallélisme suivantes ne s'appliquent pas.
Valeurs disponibles
AUTO (par défaut) : Si le parallélisme du récepteur (sink) n'est pas égal à 1 et diffère du parallélisme de l'opérateur en amont, Flink effectue automatiquement un brassage par hachage sur la clé primaire lorsque les données circulent vers le récepteur.
FORCE : Si le parallélisme du récepteur n'est pas égal à 1, Flink force un brassage par hachage sur la clé primaire lorsque les données circulent vers le récepteur.
NONE : Flink n'effectue aucun brassage par hachage basé sur le parallélisme du récepteur et de l'opérateur en amont.
Exemples
-
Définir le paramètre sur AUTO
-
Créez une tâche de streaming SQL, copiez le code SQL suivant et deploy la tâche. Le code définit explicitement le parallélisme du récepteur sur 2.
CREATE TEMPORARY TABLE s1 ( a INT, b INT, ts TIMESTAMP(3) ) WITH ( 'connector'='datagen', 'rows-per-second'='1', 'fields.ts.kind'='random','fields.ts.max-past'='5s', 'fields.b.kind'='random','fields.b.min'='0','fields.b.max'='10' ); CREATE TEMPORARY TABLE sink ( a INT, b INT, ts TIMESTAMP(3), PRIMARY KEY (a) NOT ENFORCED ) WITH ( 'connector'='print', --You can directly specify the sink parallelism by using the sink.parallelism parameter. 'sink.parallelism'='2' ); INSERT INTO sink SELECT * FROM s1; --You can also specify the sink parallelism by using dynamic table options. --INSERT INTO sink /*+ OPTIONS('sink.parallelism' = '2') */ SELECT * FROM s1; Sur la page Deployments , sous l'onglet Configuration , dans la section Resources , définissez Parallelism sur 1. Dans la section Parameters sous Other Configuration , ne définissez pas le paramètre
table.exec.sink.keyed-shuffleou ajoutez explicitementtable.exec.sink.keyed-shuffle: AUTO(les deux options ont le même effet).-
Start la tâche. Sous l'onglet Status , la connexion de données entre l'opérateur en amont et le récepteur est de type HASH.

-
-
Définir le paramètre sur FORCE
-
Créez une tâche de streaming SQL, copiez le code SQL suivant et deploy la tâche. Ce code ne spécifie pas explicitement le parallélisme du récepteur.
CREATE TEMPORARY TABLE s1 ( a INT, b INT, ts TIMESTAMP(3) ) WITH ( 'connector'='datagen', 'rows-per-second'='1', 'fields.ts.kind'='random','fields.ts.max-past'='5s', 'fields.b.kind'='random','fields.b.min'='0','fields.b.max'='10' ); CREATE TEMPORARY TABLE sink ( a INT, b INT, ts TIMESTAMP(3), PRIMARY KEY (a) NOT ENFORCED ) WITH ( 'connector'='print' ); INSERT INTO sink SELECT * FROM s1; Dans la section Resources de l'onglet Configuration de la page Deployments , définissez Parallelism sur 2. Dans la section Parameters , ajoutez
table.exec.sink.keyed-shuffle: FORCEà Other Configuration .-
Une fois la tâche start , accédez à l'onglet Status . Le parallélisme du récepteur et de l'opérateur en amont est de 2, et la connexion de données est passée en mode HASH.

-
table.exec.mini-batch.size
Ce paramètre contrôle le nombre maximal d'enregistrements pouvant être mis en mémoire tampon pour une opération par micro-lots. Lorsque ce seuil est atteint, il déclenche le calcul et émet les données. Ce paramètre prend effet uniquement lorsqu'il est utilisé conjointement avec table.exec.mini-batch.enabled et table.exec.mini-batch.allow-latency. Pour plus d'informations sur les optimisations MiniBatch, consultez MiniBatch Aggregation et MiniBatch Regular Joins .
Remarques d'utilisation
Avant le démarrage d'une tâche, si vous ne définissez pas explicitement ce paramètre dans la section Parameters, la mémoire gérée sert à mettre les données en mémoire tampon en mode mini-lot. L'une des conditions suivantes déclenche le calcul final et la sortie des données :
Réception d'un message de watermark depuis l'opérateur MiniBatchAssigner.
La mémoire gérée est pleine.
Avant le début d'un point de contrôle (checkpoint).
Arrêt de la tâche.
Valeurs disponibles
-1 (par défaut) : Indique que la mémoire gérée sert à mettre les données en mémoire tampon.
Autres valeurs Long négatives : Identique au paramètre par défaut.
Autres valeurs Long positives : Indique que la mémoire heap sert à mettre les données en mémoire tampon. Lorsque le nombre d'enregistrements en mémoire tampon atteint cette valeur (N), le système déclenche automatiquement l'opération de sortie.
Exemple
-
Créez une tâche de streaming SQL, copiez le code SQL suivant et deploy la tâche.
CREATE TEMPORARY TABLE s1 ( a INT, b INT, ts TIMESTAMP(3), PRIMARY KEY (a) NOT ENFORCED, WATERMARK FOR ts AS ts - INTERVAL '1' SECOND ) WITH ( 'connector'='datagen', 'rows-per-second'='1', 'fields.ts.kind'='random', 'fields.ts.max-past'='5s', 'fields.b.kind'='random', 'fields.b.min'='0', 'fields.b.max'='10' ); CREATE TEMPORARY TABLE sink ( a INT, b BIGINT, PRIMARY KEY (a) NOT ENFORCED ) WITH ( 'connector'='print' ); INSERT INTO sink SELECT a, sum(b) FROM s1 GROUP BY a; Sur l'onglet Configuration de la page Deployments , dans le champ Other Configuration de la section Parameters , définissez les paramètres
table.exec.mini-batch.enabled: trueettable.exec.mini-batch.allow-latency: 2s, et ne définissez pastable.exec.mini-batch.sizeafin d'utiliser sa valeur par défaut (-1).Start la tâche. Sous l'onglet Status , la topologie de la tâche inclut les opérateurs MiniBatchAssigner, LocalGroupAggregate et GlobalGroupAggregate.
table.exec.agg.mini-batch.output-identical-enabled
Lorsque State TTL is enabled , les nœuds MinibatchGlobalAgg et MinibatchAgg n'envoient pas de données en double en aval par défaut si le résultat de l'agrégation ne change pas après la consommation des données. Cela peut entraîner l'expiration de l'état des nœuds avec état en aval, car ils ne reçoivent aucune donnée de l'amont pendant une période prolongée. Ce paramètre contrôle s'il faut continuer à envoyer des données en double en aval lorsque State TTL est activé et que le résultat de l'agrégation reste inchangé. Définissez ce paramètre sur true pour que les nœuds MinibatchGlobalAgg et MinibatchAgg envoient des données dans ce cas. Si le résultat de l'agrégation de votre tâche change plus fréquemment que le State TTL configuré, vous n'avez pas besoin de définir manuellement ce paramètre. Pour plus de détails sur le problème communautaire, consultez FLINK-33936 .
Remarques d'utilisation
Ce paramètre est effectif uniquement dans VVR 8.0.8 et versions ultérieures. Dans les versions antérieures à VVR 8.0.8, le comportement équivaut à la définition de ce paramètre sur false.
Lors du changement de la valeur de false à true, la quantité de données envoyées en aval depuis les opérateurs MinibatchGlobalAgg et MinibatchAgg peut augmenter, ce qui exerce une pression accrue sur les opérateurs en aval.
Valeurs disponibles
false (par défaut) : Lorsque state TTL est activé, les opérateurs MinibatchGlobalAgg et MinibatchAgg n'émettent pas de données en aval si le résultat de l'agrégation ne change pas.
true : Lorsque state TTL est activé, les opérateurs MinibatchGlobalAgg et MinibatchAgg émettent toujours des enregistrements mis à jour (en double) en aval, même si le résultat de l'agrégation ne change pas.
Exemple
-
Créez une tâche de streaming SQL, copiez le code SQL suivant et deploy la tâche.
create temporary table src( a int, b string ) with ( 'connector' = 'datagen', 'rows-per-second' = '10', 'fields.a.min' = '1', 'fields.a.max' = '1', 'fields.b.length' = '3' ); create temporary table snk( a int, max_length_b bigint ) with ( 'connector' = 'blackhole' ); insert into snk select a, max(CHAR_LENGTH(b)) from src group by a; Dans la section Other Configuration de la zone Parameters de l'onglet Configuration de la page Deployments , définissez les paramètres
table.exec.mini-batch.enabled: trueettable.exec.mini-batch.allow-latency: 2spour activer l'optimisation Minibatch Aggregate.-
Start la tâche. Sous l'onglet Status , la tâche inclut un opérateur MinibatchGlobalAggregate. Cliquez sur le signe « + » de l'opérateur pour constater que l'opérateur GlobalGroupAggregate n'envoie pas de données en aval lorsque le résultat de l'agrégation est inchangé.
L'opérateur affiche RecordsIn à 19 et RecordsOut à 1, ce qui signifie que 19 enregistrements d'entrée n'ont produit qu'une seule sortie agrégée.
Arrêtez la tâche et ajoutez le paramètre
table.exec.agg.mini-batch.output-identical-enabled: trueà Other Configuration dans la section Parameters de la page Configuration de la page Deployments .Start la tâche. Sous l'onglet Status , vous pouvez voir que la tâche inclut un opérateur MinibatchGlobalAggregate. Cliquez sur le signe « + » de l'opérateur pour observer que l'opérateur GlobalGroupAggregate envoie désormais des données en aval, même lorsque le résultat de l'agrégation est inchangé. Après le redémarrage de la tâche, l'onglet Status indique que les valeurs RecordsIn et RecordsOut de l'opérateur GlobalGroupAggregate sont toutes deux de 94. Cela indique qu'avec
table.exec.agg.mini-batch.output-identical-enabled: trueactivé, l'opérateur envoie des données en aval même si le résultat de l'agrégation ne change pas.
table.exec.async-lookup.key-ordered-enabled
Lorsque vous utilisez une jointure de table de dimension pour l'enrichissement des données, l'activation du mode asynchrone permet souvent d'améliorer le débit. Dans une jointure de table de dimension, le paramètre table.exec.async-lookup.output-mode et le fait que l'entrée soit un flux de mises à jour déterminent l'ordre de sortie des opérations d'E/S asynchrones.
|
Output mode |
Update stream |
Non-update stream |
|
ORDERED |
ordered mode |
ordered mode |
|
ALLOW_UNORDERED |
ordered mode |
unordered mode |
Comme indiqué dans le tableau, la combinaison d'un flux de mises à jour et de ALLOW_UNORDERED garantit l'exactitude en utilisant le mode ordonné, mais cela sacrifie une partie du débit. Pour optimiser ce scénario, le paramètre table.exec.async-lookup.key-ordered-enabled a été introduit. Il équilibre la sémantique d'exactitude d'un flux de mises à jour avec les performances de débit des E/S asynchrones. Les messages d'un flux possédant la même clé de mise à jour (qui peut être considérée comme la clé primaire d'un journal des modifications) sont traités dans l'ordre où ils entrent dans l'opérateur.
Mode ordonné : Ce mode préserve l'ordre du flux. Les messages de résultat sont émis dans le même ordre que celui dans lequel les requêtes asynchrones ont été déclenchées (l'ordre dans lequel les messages entrent dans l'opérateur).
Mode non ordonné : Les messages de résultat sont émis dès que la requête asynchrone se termine. L'opérateur d'E/S asynchrone modifie l'ordre des messages dans le flux. Pour plus d'informations, consultez Asynchronous I/O | Apache Flink .
Cas d'utilisation
Utilisez cette optimisation pour préserver l'ordre de traitement par clé dans une jointure de table de dimension lorsque le flux contient peu de messages avec la même clé de mise à jour au fil du temps.
Dans un flux CDC (Change Data Capture) avec une clé primaire, vous effectuez un enrichissement des données via une jointure de table de dimension et écrivez dans un récepteur dont la clé primaire correspond à la clé primaire source. La clé de jointure pour la jointure de table de dimension est différente de la clé primaire, et la clé de jointure côté dimension est la clé primaire. Cette optimisation effectue un brassage par la clé primaire CDC, qui est dérivée en tant que clé de mise à jour. Par rapport à l'activation de l'optimisation SHUFFLE_HASH pour le même scénario, elle évite la génération d'un opérateur SinkMaterializer avant le récepteur à un parallélisme plus élevé. Cela permet d'éviter les problèmes de performance potentiels liés à cet opérateur, notamment l'état volumineux qui peut s'accumuler lors d'exécutions longues. Pour plus d'informations sur SinkUpsertMaterializer, consultez Recommendations .
La clé de jointure pour la jointure de table de dimension est différente de la clé primaire, la clé de jointure côté dimension est la clé primaire, et un opérateur de rang suit. Cette optimisation effectue un brassage par la clé primaire CDC, qui est dérivée en tant que clé de mise à jour. Par rapport à l'activation de l'optimisation SHUFFLE_HASH pour le même scénario, elle empêche UpdateFastRank de se dégrader en RetractRank. Pour savoir comment optimiser RetractRank en UpdateFastRank, consultez TopN optimization techniques .
Remarques d'utilisation
Si le flux ne possède pas de clé de mise à jour, la ligne entière est utilisée comme clé.
Le débit diminue si la même clé de mise à jour est fréquemment mise à jour sur une courte période, car les enregistrements pour une même clé de mise à jour sont traités dans un ordre strict.
Par rapport à la jointure de table de dimension asynchrone d'origine, le mode Key-Ordered introduit un état clé. L'activation ou la désactivation de ce mode affecte la compatibilité de l'état.
Cette fonctionnalité prend effet uniquement pour VVR 8.0.10 et versions ultérieures lorsque l'entrée de la jointure de table de dimension est un flux de mises à jour et que vous configurez
table.exec.async-lookup.output-mode='ALLOW_UNORDERED'ettable.exec.async-lookup.key-ordered-enabled='true'.
Valeurs disponibles
false (par défaut) : Désactive le mode Key-Ordered.
true : Active le mode Key-Ordered.
Exemple
-
L'exemple suivant utilise une jointure de table de dimension Hologres asynchrone. Créez une tâche de streaming SQL, copiez le code SQL suivant et deploy la tâche.
Pour plus d'informations sur le connecteur Hologres, consultez Hologres .
create TEMPORARY table bid_source( auction BIGINT, bidder BIGINT, price BIGINT, channel VARCHAR, url VARCHAR, dateTime TIMESTAMP(3), extra VARCHAR, proc_time as proctime(), WATERMARK FOR dateTime AS dateTime - INTERVAL '4' SECOND ) with ( 'connector' = 'kafka', -- A non-insert-only stream connector 'topic' = 'user_behavior', 'properties.bootstrap.servers' = 'localhost:9092', 'properties.group.id' = 'testGroup', 'scan.startup.mode' = 'earliest-offset', 'format' = 'csv' ); CREATE TEMPORARY TABLE users ( user_id STRING PRIMARY KEY NOT ENFORCED, -- Define the primary key user_name VARCHAR(255) NOT NULL, age INT NOT NULL ) WITH ( 'connector' = 'hologres', -- A connector that supports asynchronous lookup 'async' = 'true', 'dbname' = 'holo db name', --The name of your Hologres database 'tablename' = 'schema_name.table_name', --The name of the Hologres table to receive data 'username' = 'access id', --The AccessKey ID of your Alibaba Cloud account 'password' = 'access key', --The AccessKey Secret of your Alibaba Cloud account 'endpoint' = 'holo vpc endpoint', --The VPC endpoint of your Hologres instance ); CREATE TEMPORARY TABLE bh ( auction BIGINT, age int ) WITH ( 'connector' = 'blackhole' ); insert into bh SELECT bid_source.auction, u.age FROM bid_source JOIN users FOR SYSTEM_TIME AS OF bid_source.proc_time AS u ON bid_source.channel = u.user_id; Sur la page Deployments , sous l'onglet Configuration , dans la section Other Configuration de la zone Parameters , définissez les paramètres
table.exec.async-lookup.output-mode='ALLOW_UNORDERED'ettable.exec.async-lookup.key-ordered-enabled='true'.Start la tâche. Sous l'onglet Status , vous pouvez voir que l'attribut async de la tâche est KEY_ORDERED:true.
table.optimizer.window-join-enabled
Ce paramètre contrôle l'activation des opérations de window join . Lorsqu'il est activé, Flink optimise le plan d'exécution correspondant en tant que window join . Pour les petites fenêtres, cela réduit la surcharge d'état et améliore les performances. Par rapport à une regular join , une window join peut également éviter d'émettre des messages de mise à jour vers les opérateurs en aval, ce qui est utile pour les cas d'utilisation nécessitant des jointures avec des conditions de fenêtre temporelle réduite.
Remarques d'utilisation
Une window join impose des limitations supplémentaires sur la syntaxe SQL par rapport à une regular join , et ne prend pas en charge les flux de mises à jour.
Une window join présente une latence de sortie plus élevée qu'une regular join . La latence dépend de la taille de la fenêtre et de la vitesse d'avancement des watermarks de la source.
Lorsqu'elle est activée, une window join basée sur le temps d'événement ignore les données tardives, contrairement à une regular join .
Après modification de ce paramètre, vous ne pouvez pas reprendre l'exécution à partir d'un point de contrôle existant, car les structures d'état sous-jacentes des deux méthodes d'exécution sont incompatibles.
Valeurs disponibles
false (par défaut) : Les instructions pour une window join sont converties en une regular join pour l'exécution.
true : Active la window join . Les instructions correspondantes sont converties en une window join pour l'exécution.
Exemple
-
Créez une tâche de streaming SQL, copiez le code SQL suivant, définissez le paramètre
table.optimizer.window-join-enabledsurtrueà l'aide d'une instruction SET, puis exécutez le texte SQL pour afficher le plan d'exécution.SET 'table.optimizer.window-join-enabled' = 'true'; CREATE TEMPORARY TABLE LeftTable ( id VARCHAR, row_time TIMESTAMP_LTZ(3), num INT, WATERMARK FOR row_time as row_time - INTERVAL '5' SECONDS ) WITH ( 'connector'='datagen' ); CREATE TEMPORARY TABLE RightTable ( id VARCHAR, row_time TIMESTAMP_LTZ(3), num INT, WATERMARK FOR row_time as row_time - INTERVAL '10' SECONDS ) WITH ( 'connector'='datagen' ); EXPLAIN SELECT L.num as L_Num, L.id as L_Id, R.num as R_Num, R.id as R_Id, COALESCE(L.window_start, R.window_start) as window_start, COALESCE(L.window_end, R.window_end) as window_end FROM ( SELECT * FROM TABLE(TUMBLE(TABLE LeftTable, DESCRIPTOR(row_time), INTERVAL '5' MINUTES)) ) L JOIN ( SELECT * FROM TABLE(TUMBLE(TABLE RightTable, DESCRIPTOR(row_time), INTERVAL '5' MINUTES)) ) R ON L.num = R.num AND L.window_start = R.window_start AND L.window_end = R.window_end;La section Optimized Execution Plan de la sortie indique que le plan contient un opérateur WindowJoin.
== Optimized Execution Plan == Calc(select=[num AS L_Num, id AS L_Id, num0 AS R_Num, id0 AS R_Id, CASE(window_start IS NOT NULL, window_start, window_start0) AS window_start, CASE(window_end IS NOT NULL, window_end, window_end0) AS window_end]) +- WindowJoin(leftWindow=[TUMBLE(win_start=[window_start], win_end=[window_end], size=[5 min])], rightWindow=[TUMBLE(win_start=[window_start], win_end=[window_end], size=[5 min])], joinType=[InnerJoin], where=[(num = num0)], select=[id, num, window_start, window_end, id0, num0, window_start0, window_end0]) :- Exchange(distribution=[hash[num]]) : +- Calc(select=[id, num, window_start, window_end]) : +- WindowTableFunction(window=[TUMBLE(time_col=[row_time], size=[5 min])]) : +- WatermarkAssigner(rowtime=[row_time], watermark=[(row_time - 5000:INTERVAL SECOND)]) : +- TableSourceScan(table=[[vvp, default, LeftTable]], fields=[id, row_time, num]) +- Exchange(distribution=[hash[num]]) +- Calc(select=[id, num, window_start, window_end]) +- WindowTableFunction(window=[TUMBLE(time_col=[row_time], size=[5 min])]) +- WatermarkAssigner(rowtime=[row_time], watermark=[(row_time - 10000:INTERVAL SECOND)]) +- TableSourceScan(table=[[vvp, default, RightTable]], fields=[id, row_time, num]) -
Modifiez l'instruction SET dans le code SQL pour définir le paramètre
table.optimizer.window-join-enabledsurfalseou supprimez l'instruction SET, puis exécutez le texte SQL pour afficher le plan d'exécution modifié.-- set to 'false' or remove this setting clause SET 'table.optimizer.window-join-enabled' = 'false'; CREATE TEMPORARY TABLE LeftTable ( id VARCHAR, row_time TIMESTAMP_LTZ(3), num INT, WATERMARK FOR row_time as row_time - INTERVAL '5' SECONDS ) WITH ( 'connector'='datagen' ); CREATE TEMPORARY TABLE RightTable ( id VARCHAR, row_time TIMESTAMP_LTZ(3), num INT, WATERMARK FOR row_time as row_time - INTERVAL '10' SECONDS ) WITH ( 'connector'='datagen' ); EXPLAIN SELECT L.num as L_Num, L.id as L_Id, R.num as R_Num, R.id as R_Id, COALESCE(L.window_start, R.window_start) as window_start, COALESCE(L.window_end, R.window_end) as window_end FROM ( SELECT * FROM TABLE(TUMBLE(TABLE LeftTable, DESCRIPTOR(row_time), INTERVAL '5' MINUTES)) ) L JOIN ( SELECT * FROM TABLE(TUMBLE(TABLE RightTable, DESCRIPTOR(row_time), INTERVAL '5' MINUTES)) ) R ON L.num = R.num AND L.window_start = R.window_start AND L.window_end = R.window_end;La section Optimized Execution Plan de la sortie ne contient plus d'opérateur WindowJoin. L'opération est désormais une jointure régulière.
== Optimized Execution Plan == Calc(select=[num AS L_Num, id AS L_Id, num0 AS R_Num, id0 AS R_Id, CASE(window_start IS NOT NULL, window_start, window_start0) AS window_start, CASE(window_end IS NOT NULL, window_end, window_end0) AS window_end]) +- Join(joinType=[InnerJoin], where=[((num = num0) AND (window_start = window_start0) AND (window_end = window_end0))], select=[id, num, window_start, window_end, id0, num0, window_start0, window_end0], leftInputSpec=[NoUniqueKey], rightInputSpec=[NoUniqueKey]) :- Exchange(distribution=[hash[num, window_start, window_end]]) : +- Calc(select=[id, num, window_start, window_end]) : +- WindowTableFunction(window=[TUMBLE(time_col=[row_time], size=[5 min])]) : +- WatermarkAssigner(rowtime=[row_time], watermark=[(row_time - 5000:INTERVAL SECOND)]) : +- TableSourceScan(table=[[vvp, default, LeftTable]], fields=[id, row_time, num]) +- Exchange(distribution=[hash[num, window_start, window_end]]) +- Calc(select=[id, num, window_start, window_end]) +- WindowTableFunction(window=[TUMBLE(time_col=[row_time], size=[5 min])]) +- WatermarkAssigner(rowtime=[row_time], watermark=[(row_time - 10000:INTERVAL SECOND)]) +- TableSourceScan(table=[[vvp, default, RightTable]], fields=[id, row_time, num])