Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Développement et débogage

Dernière mise à jour :Aug 09, 2026

Cette rubrique aborde les problèmes courants liés au développement et au débogage.

Déclaration d'instructions DDL avec des instructions DML

Lorsque vous soumettez des instructions DDL et DML conjointement dans le même script, déclarez l'instruction DDL en utilisant CREATE TEMPORARY TABLE plutôt que CREATE TABLE. À défaut, la validation via Validate échouera avec une erreur semblable à celle ci-dessous.

CREATE TABLE datagen_source (a bigint, b int, c varchar)
  WITH ('connector' = 'datagen');
CREATE  TABLE print_sink(C bigint, var1 int)
  WITH ('connector' = 'print','logger' = 'true');
INSERT INTO print_sink SELECT a,8  FROM datagen_source;
Error message:
org.apache.flink.table.gateway.api.vvr.utils.SqlValidationException: A sequence of multiple statements to execute is supported if the last statement is a 'SELECT' statement or 'INSERT INTO' statement or 'CREATE TABLE IF NOT EXISTS ... AS TABLE' statement ...BASE IF NOT EXISTS ... AS DATABASE' statement or 'AUTO OPTIMIZE TABLE|DATABASE' statements or multiple 'INSERT INTO' or 'CREATE TABLE IF NOT EXISTS ... AS TAB ...ATABASE IF NOT EXISTS ... AS DATABASE' statements wrapped in a 'BEGIN STATEMENT SET' block and all other statements are CREATE TEMPORARY TABLE|VIEW|[SYSTEM] FUNCTION, 'SHOW', DESCRIBE, 'USE' statements.
	at org.apache.flink.table.gateway.vvr.service.utils.SqlValidateUtils.validateDraft(SqlValidateUtils.java:107)
	at org.apache.flink.table.gateway.vvr.service.command.DraftCommand.getDraftType(DraftCommand.java:120)
	at org.apache.flink.table.gateway.vvr.service.command.DraftCommand.executeInternal(DraftCommand.java:71)
	at java.security.AccessController.doPrivileged(Native Method)
	at javax.security.auth.Subject.doAs(Subject.java:422)

Instructions INSERT INTO multiples

Pour former une seule unité logique, encapsulez plusieurs instructions INSERT INTO entre BEGIN STATEMENT SET; et END;. Pour plus d'informations, consultez la rubrique Instructions INSERT INTO. Si vous n'encapsulez pas les instructions, la validation via Validate échouera avec une erreur semblable à celle ci-dessous.

CREATE TEMPORARY TABLE datagen_source (a bigint, b int, c varchar)
  WITH ('connector' = 'datagen');
CREATE TEMPORARY  TABLE print_sink(C bigint, var1 int)
  WITH ('connector' = 'print','logger' = 'true');
CREATE TEMPORARY TABLE print_sink2(C bigint, var2 int)
  WITH ('connector' = 'print','logger' = 'true');
INSERT INTO print_sink SELECT a,B  FROM datagen_source;
INSERT INTO print_sink2 SELECT a,B  FROM datagen_source;
org.apache.flink.table.gateway.api.vvr.utils.SqlValidationException: A sequence of multiple statements to execute is supported if the last statement is a 'SELECT' statement or 'INSERT INTO' statement or 'CREATE TABLE IF NOT EXISTS ... AS TABLE' statement or 'CREATE DATABASE IF NOT EXISTS ... AS DATABASE' statement or 'AUTO OPTIMIZE TABLE|DATABASE' statements or multiple 'INSERT INTO' or 'CREATE TABLE IF NOT EXISTS ... AS TABLE' or 'CREATE DATABASE IF NOT EXISTS ... AS DATABASE' statements wrapped in a 'BEGIN STATEMENT SET' block and all other statements are CREATE TEMPORARY TABLE|VIEW|[SYSTEM] FUNCTION, 'SHOW', DESCRIBE, 'USE' statements.
      at org.apache.flink.table.gateway.vvr.service.utils.SqlValidateUtils.validateDraft(SqlValidateUtils.java:79)
      at org.apache.flink.table.gateway.vvr.service.command.DraftCommand.getDraftType(DraftCommand.java:120)
      at org.apache.flink.table.gateway.vvr.service.command.DraftCommand.executeInternal(DraftCommand.java:71)
      at java.security.AccessController.doPrivileged(Native Method)
      at javax.security.auth.Subject.doAs(Subject.java:422)
      at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1899)
      at org.apache.flink.table.gateway.service.context.SqlGatewaySecurityContext.runSecured(SqlGatewaySecurityContext.java:73)
      at org.apache.flink.table.gateway.vvr.service.command.AbstractCommand.wrapClassLoader(AbstractCommand.java:171)
      at org.apache.flink.table.gateway.vvr.service.command.AbstractCommand.execute(AbstractCommand.java:163)
      at org.apache.flink.table.gateway.vvr.service.command.CommandManager.lambda$execute$0(CommandManager.java:71)
      at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
      at java.util.concurrent.FutureTask.run(FutureTask.java:266)
      at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)

Transmission de caractères spéciaux dans les arguments principaux

  • Cause

    Lorsque vous transmettez des caractères spéciaux tels que # et $ dans les Entry point Main Arguments, le caractère d'échappement backslash (\) ne fonctionne pas et les caractères spéciaux sont ignorés.

  • Solution

    Sur la page Deployments, cliquez sur le nom du déploiement cible. Dans la section Parameters, ajoutez le paramètre env.java.opts: -Dconfig.disable-inline-comment=true au champ Other Configuration. Pour plus d'informations, consultez la rubrique Configuration des paramètres de déploiement personnalisés.

Échec du téléchargement du JAR UDF après modification

  • Cause

    L'environnement d'exécution UDF impose l'unicité des noms de classes entre les packages JAR.

  • Solution

    • Supprimez l'ancien package JAR et téléchargez le nouveau.

    • Téléchargez le package JAR dans la section Additional Dependencies et utilisez une fonction temporaire dans votre code. Pour savoir comment utiliser une fonction temporaire, consultez la rubrique Enregistrement d'une UDF. L'exemple suivant illustre la syntaxe.

      CREATE TEMPORARY FUNCTION `cp_record_reduce` AS 'com.taobao.test.udf.blink.CPRecordReduceUDF';

      Dans la section Additional Dependencies, située dans le panneau Other Configuration à droite, indiquez l'URL OSS de votre package JAR UDF.

Désalignement des champs avec une classe POJO comme type de retour UDTF****

  • Détails

    Lorsque vous utilisez une classe POJO comme type de retour pour une UDTF et déclarez explicitement une liste d'alias pour les colonnes renvoyées en SQL, vous pouvez rencontrer un problème de désalignement des champs. Par conséquent, les champs réels peuvent ne pas correspondre à ceux attendus, même si les types de données sont cohérents.

    Par exemple, si vous utilisez la classe POJO suivante comme type de retour pour une UDTF, que vous l'emballiez et que vous enregistriez la fonction en tant qu'UDF au niveau du déploiement comme décrit dans la rubrique Développement d'une fonction personnalisée, la validation SQL échoue.

    package com.aliyun.example;
    public class TestPojoWithoutConstructor {
    	public int c;
    	public String d;
    	public boolean a;
    	public String b;
    }
    package com.aliyun.example;
    import org.apache.flink.table.functions.TableFunction;
    public class MyTableFuncPojoWithoutConstructor extends TableFunction<TestPojoWithoutConstructor> {
    	private static final long serialVersionUID = 1L;
    	public void eval(String str1, Integer i2) {
    		TestPojoWithoutConstructor p = new TestPojoWithoutConstructor();
    		p.d = str1 + "_d";
    		p.c = i2 + 2;
    		p.b = str1 + "_b";
    		collect(p);
    	}
    }
    CREATE TEMPORARY FUNCTION MyTableFuncPojoWithoutConstructor as 'com.aliyun.example.MyTableFuncPojoWithoutConstructor';
    CREATE TEMPORARY TABLE src ( 
      id STRING,
      cnt INT
    ) WITH (
      'connector' = 'datagen'
    );
    CREATE TEMPORARY TABLE sink ( 
      f1 INT,
      f2 STRING,
      f3 BOOLEAN,
      f4 STRING
    ) WITH (
     'connector' = 'print'
    );
    INSERT INTO sink
    SELECT T.* FROM src, LATERAL TABLE(MyTableFuncPojoWithoutConstructor(id, cnt)) AS T(c, d, a, b);

    La validation SQL renvoie le message d'erreur suivant :

    org.apache.flink.table.api.ValidationException: SQL validation failed. Column types of query result and sink for 'vvp.default.sink' do not match.
    Cause: Sink column 'f1' at position 0 is of type INT but expression in the query is of type BOOLEAN NOT NULL.
    Hint: You will need to rewrite or cast the expression.
    Query schema: [c: BOOLEAN NOT NULL, d: STRING, a: INT NOT NULL, b: STRING]
    Sink schema:  [f1: INT, f2: STRING, f3: BOOLEAN, f4: STRING]
    	at org.apache.flink.table.sqlserver.utils.FormatValidatorExceptionUtils.newValidationException(FormatValidatorExceptionUtils.java:41)

    Les champs renvoyés par l'UDTF sont désalignés par rapport aux champs de la classe POJO. Dans le résultat de la requête, le champ c est de type BOOLEAN et le champ a est de type INT, ce qui est l'inverse de leurs définitions dans la classe POJO.

  • Cause

    Selon les règles d'inférence de type pour les classes POJO :

    • Si la classe POJO possède un constructeur paramétré, Flink déduit le type de retour en se basant sur l'ordre des paramètres du constructeur.

    • Si la classe POJO ne possède pas de constructeur paramétré, Flink réorganise les champs par ordre alphabétique selon leur nom.

    Dans l'exemple, comme la classe POJO utilisée pour le type de retour de l'UDTF ne possède pas de constructeur paramétré, les champs sont renvoyés par ordre alphabétique, résultant en le type BOOLEAN a, VARCHAR(2147483647) b, INTEGER c, VARCHAR(2147483647) d). Bien que cette inférence soit valide, la requête SQL renomme explicitement les colonnes de sortie avec LATERAL TABLE(MyTableFuncPojoWithoutConstructor(id, cnt)) AS T(c, d, a, b). Cette liste d'alias renomme les colonnes par position, créant une incompatibilité avec les champs triés alphabétiquement issus du POJO. Ce conflit entre l'aliasing positionnel et l'ordre alphabétique des champs provoque l'exception de validation ou un désalignement inattendu des données.

  • Solution

    • Si la classe POJO ne possède pas de constructeur paramétré, supprimez le renommage explicite des champs de retour de l'UDTF. Par exemple, modifiez l'instruction INSERT dans le SQL comme suit :

      -- If the POJO class lacks a parameterized constructor, select the required fields by name. 
      -- When using T.*, you must be aware of the actual order of the returned fields.
      SELECT T.c, T.d, T.a, T.b FROM src, LATERAL TABLE(MyTableFuncPojoWithoutConstructor(id, cnt)) AS T;
    • Implémentez un constructeur paramétré dans la classe POJO pour contrôler l'ordre des champs dans le type de retour. Dans ce cas, l'ordre des champs de la sortie de l'UDTF correspondra à l'ordre des paramètres du constructeur.

      package com.aliyun.example;
      public class TestPojoWithConstructor {
      	public int c;
      	public String d;
      	public boolean a;
      	public String b;
      	// Using specific fields order instead of alphabetical order
      	public TestPojoWithConstructor(int c, String d, boolean a, String b) {
      		this.c = c;
      		this.d = d;
      		this.a = a;
      		this.b = b;
      	}
      }

Résolution des conflits de dépendances Flink

  • Symptômes

    • Le conflit se manifeste par des erreurs claires levées par des classes liées à Flink ou Hadoop.

      java.lang.AbstractMethodError
      java.lang.ClassNotFoundException
      java.lang.IllegalAccessError
      java.lang.IllegalAccessException
      java.lang.InstantiationError
      java.lang.InstantiationException
      java.lang.InvocationTargetException
      java.lang.NoClassDefFoundError
      java.lang.NoSuchFieldError
      java.lang.NoSuchFieldException
      java.lang.NoSuchMethodError
      java.lang.NoSuchMethodException
    • Alternativement, le problème peut se présenter sous la forme d'un comportement inattendu sans message d'erreur clair, tel que :

      • Les journaux ne sont pas générés ou la configuration log4j ne prend pas effet.

        Ce problème est généralement causé par des configurations liées à log4j incluses dans les dépendances. Vérifiez si le package JAR de déploiement contient des dépendances qui embarquent des configurations log4j. Vous pouvez supprimer ces configurations en utilisant des exclusions dans vos définitions de dépendances.

        Remarque

        Si vous devez utiliser une version différente de log4j, utilisez le plugin maven-shade-plugin pour relocaliser les classes liées à log4j.

      • Exceptions d'appel RPC.

        Les conflits de dépendances affectant les appels RPC Akka de Flink peuvent provoquer des exceptions qui ne s'affichent pas dans les journaux par défaut. Vous devez activer la journalisation de débogage pour les identifier.

        Par exemple, le journal de débogage affiche Cannot allocate the requested resources. Trying to allocate ResourceProfile{xxx}, mais le journal du JobManager (JM) ne montre aucune activité après Registering TaskManager with ResourceID xxx jusqu'à ce qu'une erreur de délai d'attente NoResourceAvailableException se produise. Pendant ce temps, le TaskManager (TM) signale continuellement l'erreur Cannot allocate the requested resources. Trying to allocate ResourceProfile{xxx}.

        Cause : Avec la journalisation de débogage activée, vous pouvez voir qu'une InvocationTargetException est levée lors d'un appel RPC. Cette erreur entraîne l'échec de l'allocation de slot TM en cours de processus, ce qui provoque un état incohérent. Le ResourceManager (RM) tente alors continuellement et sans succès d'allouer un slot, sans pouvoir récupérer.

  • Causes

    • Le package JAR de déploiement contient des dépendances inutiles, telles que les bibliothèques de base Flink, Hadoop ou log4j, qui provoquent des conflits de dépendances.

    • Les dépendances pour un connecteur requis ne sont pas incluses dans le package JAR.

  • Dépannage

    • Examinez le fichier pom.xml du déploiement pour identifier les dépendances inutiles.

    • Inspectez le contenu du package JAR de déploiement en exécutant la commande jar tf foo.jar pour vérifier la présence de fichiers conflictuels.

    • Analysez l'arborescence des dépendances du déploiement pour détecter les conflits en exécutant la commande mvn dependency:tree.

  • Solution

    • En bonne pratique, définissez la valeur scope des dépendances du framework de base sur provided. Cela empêche leur inclusion dans le package JAR de déploiement.

      • DataStream Java

        <dependency>
          <groupId>org.apache.flink</groupId>
          <artifactId>flink-streaming-java_2.11</artifactId>
          <version>${flink.version}</version>
          <scope>provided</scope>
        </dependency>
      • DataStream Scala

        <dependency>
          <groupId>org.apache.flink</groupId>
          <artifactId>flink-streaming-scala_2.11</artifactId>
          <version>${flink.version}</version>
          <scope>provided</scope>
        </dependency>
      • DataSet Java

        <dependency>
          <groupId>org.apache.flink</groupId>
          <artifactId>flink-java</artifactId>
          <version>${flink.version}</version>
          <scope>provided</scope>
        </dependency>
      • DataSet Scala

        <dependency>
          <groupId>org.apache.flink</groupId>
          <artifactId>flink-scala_2.11</artifactId>
          <version>${flink.version}</version>
          <scope>provided</scope>
        </dependency>
    • Ajoutez les dépendances de connecteur requises à votre projet. La portée par défaut est compile, ce qui permet de les inclure correctement dans le package JAR de déploiement. Par exemple, pour ajouter le connecteur Kafka :

      <dependency>
          <groupId>org.apache.flink</groupId>
          <artifactId>flink-connector-kafka_2.11</artifactId>
          <version>${flink.version}</version>
      </dependency>
    • N'ajoutez pas d'autres dépendances Flink, Hadoop ou log4j. Toutefois :

      • Si le déploiement possède une dépendance directe sur des composants de configuration de base ou liés aux connecteurs, définissez la portée sur provided. L'exemple suivant illustre la syntaxe.

        <dependency>
            <groupId>org.apache.hadoop</groupId>
            <artifactId>hadoop-common</artifactId>
            <scope>provided</scope>
        </dependency>
      • Si le déploiement possède une dépendance transitive sur des composants de configuration de base ou liés aux connecteurs, supprimez-la en utilisant une exclusion. L'exemple suivant illustre la syntaxe.

        <dependency>
            <groupId>foo</groupId>
              <artifactId>bar</artifactId>
              <exclusions>
                <exclusion>
                <groupId>org.apache.hadoop</groupId>
                <artifactId>hadoop-common</artifactId>
               </exclusion>
            </exclusions>
        </dependency>

Erreur : Could not parse type at position 50: expected but was . Input type string: ROW

  • Détails de l'erreur

    Lors de l'écriture de SQL dans l'éditeur SQL, une erreur de vérification de syntaxe (ligne ondulée rouge) se produit lorsque vous utilisez une UDTF.

    Caused by: org.apache.flink.table.api.ValidationException: Could not parse type at position 50: <IDENTIFIER> expected but was <KEYWORD>. Input type string: ROW<resultId String,pointRange String,from String,to String,type String,pointScope String,userId String,point String,triggerSource String,time String,uuid String>

    Le code suivant est un exemple :

    @FunctionHint(
        //input = @DataTypeHint("BYTES"),
        output = @DataTypeHint("ROW<resultId String,pointRange String,from String,to String,type String,pointScope String,userId String,point String,triggerSource String,time String,uuid String>"))
    public class PointChangeMetaQPaser1 extends TableFunction<Row> {
        Logger logger = LoggerFactory.getLogger(this.getClass().getName());
        public void eval(byte[] bytes) {
            try {
                String messageBody = new String(bytes, "UTF-8");
                Map<String, String> resultDO = JSON.parseObject(messageBody, Map.class);
                logger.info("PointChangeMetaQPaser1 logger:" + JSON.toJSONString(resultDO));
                collect(Row.of(
                        getString(resultDO.get("resultId")),
                        getString(resultDO.get("pointRange")),
                        getString(resultDO.get("from")),
                        getString(resultDO.get("to")),
                        getString(resultDO.get("type")),
                        getString(resultDO.get("pointScope")),
                        getString(resultDO.get("userId")),
                        getString(resultDO.get("point")),
                        getString(resultDO.getOrDefault("triggerSource", "NULL")),
                        getString(resultDO.getOrDefault("time", String.valueOf(System.currentTimeMillis()))),
                        getString(resultDO.getOrDefault("uuid", String.valueOf(UUID.randomUUID())))
                ));
            } catch (Exception e) {
                logger.error("PointChangeMetaQPaser1 error", e);
            }
        }
        private String getString(Object o) {
            if (o == null) {
                return null;
            }
            return String.valueOf(o);
        }
    }
  • Cause

    Lorsque vous utilisez DataTypeHint pour définir les types de données de fonction, un mot-clé réservé est utilisé directement comme nom de champ.

  • Solution

    • Renommez les champs en évitant les mots-clés. Par exemple, renommez to en fto et from en ffrom.

    • Entourez les champs utilisant des mots-clés réservés de backticks (`).

Erreur : « Invalid primary key. Column 'xxx' is nullable. »

  • Cause

    Flink impose que toutes les colonnes de clé primaire soient explicitement déclarées comme NOT NULL. Même si les données ne contiennent aucune valeur NULL, Flink rejette l'opération avant l'écriture si une colonne de clé primaire dans l'instruction de création de table autorise les valeurs NULL (par exemple, INT NULL). Il ne s'agit pas d'une erreur d'exécution, mais d'une vérification sémantique lors de la phase d'analyse DDL.

  • Solution

    Déclarez les colonnes de clé primaire mentionnées dans l'erreur comme NOT NULL et recréez la table.

Ouverture d'un fichier JSON dans le navigateur au lieu du téléchargement

  • Symptôme

    Lorsque vous cliquez pour télécharger un fichier JSON depuis la page Artifacts, le navigateur ne déclenche pas le téléchargement. Au lieu de cela, il ouvre un nouvel onglet et affiche directement le contenu JSON.

  • Cause

    Le fichier JSON dans OSS est dépourvu de l'en-tête de réponse HTTP Content-Disposition: attachment. Cela amène le navigateur à afficher directement le contenu du fichier au lieu de le télécharger.

  • Solution

    • Option 1 : Retélécharger le fichier

      Ce problème a été corrigé dans la version 4.5.0 de la plateforme, mais le correctif ne s'applique qu'aux fichiers téléchargés après la publication de cette version. Les fichiers téléchargés avant cette date doivent être traités manuellement.

    • Option 2 : Modifier les métadonnées de l'objet OSS

      Modifiez manuellement les métadonnées de l'objet en ajoutant l'attribut HTTP standard suivant :

      • Nom de l'en-tête : Content-Disposition

      • Valeur de l'en-tête : attachment

      Pour plus d'informations, consultez la rubrique Gestion des métadonnées d'objet.