Todos os produtos
Search
Central de documentação

MaxCompute:Gravar dados no MaxCompute com Flink

Última atualização: Jun 27, 2026

Use o conector Flink do MaxCompute para gravar dados do Flink em tabelas padrão e delta tables no MaxCompute e simplificar a ingestão de dados. Este tópico descreve os recursos do conector e detalha os procedimentos de gravação.

Contexto

  • Modos de gravação compatíveis

    O conector Flink oferece dois modos de gravação: upsert e insert. No modo upsert, agrupe os fluxos de dados das seguintes formas:

    • Agrupamento por chave primária

    • Agrupamento por campo de partição

      Embora adequado para um grande número de partições, o agrupamento por campo de partição pode causar skew de dados.

  • Para obter detalhes sobre o procedimento de gravação upsert e os parâmetros recomendados, consulte Ingestão de dados em tempo real em um data warehouse.

  • Especifique o modo de gravação por meio dos parâmetros do conector Flink. Para a lista completa de parâmetros, consulte Apêndice: Parâmetros do conector Flink.

  • Defina o intervalo de checkpoint para jobs de gravação upsert do Flink em pelo menos 3 minutos. Intervalos menores reduzem a eficiência da gravação e geram muitos arquivos pequenos.

  • A tabela a seguir apresenta o mapeamento de tipos de dados entre Realtime Compute for Apache Flink e MaxCompute.

    Tipo de dado Flink

    Tipo de dado MaxCompute

    CHAR(p)

    CHAR(p)

    VARCHAR(p)

    VARCHAR(p)

    STRING

    STRING

    BOOLEAN

    BOOLEAN

    TINYINT

    TINYINT

    SMALLINT

    SMALLINT

    INT

    INT

    BIGINT

    BIGINT

    FLOAT

    FLOAT

    DOUBLE

    DOUBLE

    DECIMAL(p, s)

    DECIMAL(p, s)

    DATE

    DATE

    TIMESTAMP(9) WITHOUT TIME ZONE, TIMESTAMP_LTZ(9)

    TIMESTAMP

    TIMESTAMP(3) WITHOUT TIME ZONE, TIMESTAMP_LTZ(3)

    DATETIME

    BYTES

    BINARY

    ARRAY<T>

    ARRAY<T>

    MAP<K, V>

    MAP<K, V>

    ROW

    STRUCT

    Nota

    O tipo de dado TIMESTAMP do Flink não inclui informações de fuso horário, enquanto o TIMESTAMP do MaxCompute inclui. Essa diferença pode causar uma discrepância de 8 horas. Para alinhar os timestamps, use TIMESTAMP_LTZ(9).

    -- Flink SQL
    CREATE TEMPORARY TABLE odps_source(
      id BIGINT NOT NULL COMMENT 'ID',
      created_time TIMESTAMP NOT NULL COMMENT 'Creation time',
      updated_time TIMESTAMP_LTZ(9) NOT NULL COMMENT 'Update time',
    PRIMARY KEY (id) NOT ENFORCED
    ) WITH (
    'connector' = 'maxcompute',
    ...
    );

Gravar dados de um cluster Flink autogerenciado

  1. Pré-requisitos: Crie uma tabela no MaxCompute.

    Crie primeiro uma tabela no MaxCompute para receber os dados do Flink. O exemplo abaixo demonstra esse processo com a criação de duas tabelas (uma delta table não particionada e uma tabela particionada). Para mais informações sobre as configurações de propriedades da tabela, consulte parâmetros de delta table.

    -- Create a non-partitioned delta table.
    CREATE TABLE mf_flink_tt (
      id BIGINT not null,
      name STRING,
      age INT,
      status BOOLEAN, primary key (id)
    )
    tblproperties ("transactional"="true", 
                   "write.bucket.num" = "64", 
                   "acid.data.retain.hours"="12") ;
    
    --Create a partitioned delta table.
    CREATE TABLE mf_flink_tt_part (
      id BIGINT not null,
      name STRING,
      age INT,
      status BOOLEAN, 
      primary key (id)
    )
      partitioned by (dd string, hh string) 
      tblproperties ("transactional"="true", 
                     "write.bucket.num" = "64", 
                     "acid.data.retain.hours"="12") ;
    
  2. Configure um cluster Flink open source. O conector é compatível com as versões 1,13, 1,15, 1,16 e 1,17 do Flink. Selecione o conector Flink correspondente à sua versão:

    Nota
    • O conector Flink para Flink 1.16 é compatível com Flink 1.17.

    • Este tópico usa o conector Flink para Flink 1.13 como exemplo. Baixe e descompacte o pacote.

  3. Baixe o conector Flink e adicione-o ao pacote do seu cluster Flink.

    1. Baixe o pacote JAR do conector Flink para o seu ambiente local.

    2. Adicione o pacote JAR do conector Flink ao diretório lib do pacote de instalação descompactado do Flink.

      mv flink-connector-odps-1.13-shaded.jar $FLINK_HOME/lib/flink-connector-odps-1.13-shaded.jar
  4. Inicie o serviço da instância Flink.

    cd $FLINK_HOME/bin
    ./start-cluster.sh
  5. Inicie o cliente Flink.

    cd $FLINK_HOME/bin
    ./sql-client.sh
  6. Crie uma tabela Flink e configure os parâmetros do conector.

    Crie uma tabela Flink e configure seus parâmetros usando Flink SQL ou a API DataStream. As seções a seguir fornecem exemplos essenciais para ambas as abordagens.

    Flink SQL

    1. No editor Flink SQL, execute os comandos abaixo para criar uma tabela e configurar os parâmetros.

      -- Register a non-partitioned table in Flink SQL.
      CREATE TABLE mf_flink (
        id BIGINT,
        name STRING,
        age INT,
        status BOOLEAN,
        PRIMARY KEY(id) NOT ENFORCED
      ) WITH (
        'connector' = 'maxcompute',
        'table.name' = 'mf_flink_tt',
        'sink.operation' = 'upsert',
        'odps.access.id'='LTAI****************',
        'odps.access.key'='********************',
        'odps.end.point'='http://service.cn-beijing.maxcompute.aliyun.com/api',
        'odps.project.name'='mf_mc_bj'
      );
      
      -- Register a partitioned table in Flink SQL.
      CREATE TABLE mf_flink_part (
        id BIGINT,
        name STRING,
        age INT,
        status BOOLEAN,
        dd STRING,
        hh STRING,
        PRIMARY KEY(id) NOT ENFORCED
      ) PARTITIONED BY (`dd`,`hh`)
      WITH (
        'connector' = 'maxcompute',
        'table.name' = 'mf_flink_tt_part',
        'sink.operation' = 'upsert',
        'odps.access.id'='LTAI****************',
        'odps.access.key'='********************',
        'odps.end.point'='http://service.cn-beijing.maxcompute.aliyun.com/api',
        'odps.project.name'='mf_mc_bj'
      );
    2. Grave dados na tabela Flink e consulte a tabela MaxCompute para verificar o resultado.

      -- Insert data into the non-partitioned table in the Flink SQL client.
      INSERT INTO mf_flink VALUES (1,'Danny',27, false);
      
      -- Query result in MaxCompute.
      SELECT * FROM mf_flink_tt;
      +------------+------+------+--------+
      | id         | name | age  | status |
      +------------+------+------+--------+
      | 1          | Danny | 27   | false  |
      +------------+------+------+--------+
      
      -- Insert data into the non-partitioned table in the Flink SQL client to update the record.
      INSERT INTO mf_flink VALUES (1,'Danny',28, false);
      -- Query result in MaxCompute.
      SELECT * FROM mf_flink_tt;
      +------------+------+------+--------+
      | id         | name | age  | status |
      +------------+------+------+--------+
      | 1          | Danny | 28   | false  |
      +------------+------+------+--------+
      
      -- Insert data into the partitioned table in the Flink SQL client.
      INSERT INTO mf_flink_part VALUES (1,'Danny',27, false, '01','01');
      -- Query result in MaxCompute.
      SELECT * FROM mf_flink_tt_part WHERE dd=01 AND hh=01;
      +------------+------+------+--------+----+----+
      | id         | name | age  | status | dd | hh |
      +------------+------+------+--------+----+----+
      | 1          | Danny | 27   | false  | 01 | 01 |
      +------------+------+------+--------+----+----+
      
      -- Insert data into the partitioned table in the Flink SQL client to update the record.
      INSERT INTO mf_flink_part VALUES (1,'Danny',30, false, '01','01');
      -- Query result in MaxCompute.
      SELECT * FROM mf_flink_tt_part WHERE dd=01 AND hh=01;
      +------------+------+------+--------+----+----+
      | id         | name | age  | status | dd | hh |
      +------------+------+------+--------+----+----+
      | 1          | Danny | 30   | false  | 01 | 01 |
      +------------+------+------+--------+----+----+

    DataStream API

    1. Para usar a API DataStream, adicione a seguinte dependência.

      <dependency>
        <groupId>com.aliyun.odps</groupId>
        <artifactId>flink-connector-maxcompute</artifactId>
                  <version>xxx</version>
                  <scope>system</scope>
                  <systemPath>${mvn_project.basedir}/lib/flink-connector-maxcompute-xxx-shaded.jar</systemPath>
      </dependency>
      Nota

      Substitua "xxx" pelo número real da versão.

    2. O código de exemplo a seguir mostra como criar uma tabela e configurar parâmetros.

      package com.aliyun.odps.flink.examples;
      
      import org.apache.flink.configuration.Configuration;
      import org.apache.flink.odps.table.OdpsOptions;
      import org.apache.flink.odps.util.OdpsConf;
      import org.apache.flink.odps.util.OdpsPipeline;
      import org.apache.flink.streaming.api.datastream.DataStream;
      import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
      import org.apache.flink.table.api.Table;
      import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
      import org.apache.flink.table.data.RowData;
      
      public class Examples {
      
          public static void main(String[] args) throws Exception {
              StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
              env.enableCheckpointing(120 * 1000);
      
              StreamTableEnvironment streamTableEnvironment = StreamTableEnvironment.create(env);
      
              Table source = streamTableEnvironment.sqlQuery("SELECT * FROM source_table");
              DataStream<RowData> input = streamTableEnvironment.toAppendStream(source, RowData.class);
      
              Configuration config = new Configuration();
              config.set(OdpsOptions.SINK_OPERATION, "upsert");
              config.set(OdpsOptions.UPSERT_COMMIT_THREAD_NUM, 8);
              config.set(OdpsOptions.UPSERT_MAJOR_COMPACT_MIN_COMMITS, 100);
      
              OdpsConf odpsConfig = new OdpsConf("accessid",
                      "accesskey",
                      "endpoint",
                      "project",
                      "tunnel endpoint");
      
              OdpsPipeline.Builder builder = OdpsPipeline.builder();
              builder.projectName("sql2_isolation_2a")
                      .tableName("user_ledger_portfolio")
                      .partition("")
                      .configuration(config)
                      .odpsConf(odpsConfig)
                      .sink(input, false);
              env.execute();
          }
      }

Gravar dados do Flink totalmente gerenciado

  1. Pré-requisitos: Crie uma tabela no MaxCompute.

    Crie uma tabela de destino no MaxCompute para os dados do Flink. O exemplo a seguir mostra como criar uma delta table.

    SET odps.sql.type.system.odps2=true;
    DROP TABLE mf_flink_upsert;
    CREATE TABLE mf_flink_upsert (
      c1 int not null, 
      c2 string, 
      gt timestamp,
      primary key (c1)
    ) 
      PARTITIONED BY (ds string)
      tblproperties ("transactional"="true",
                     "write.bucket.num" = "64", 
                     "acid.data.retain.hours"="12") ;
  2. O conector Flink vem pré-carregado no Flink totalmente gerenciado, portanto, nenhuma instalação manual é necessária. Visualize os detalhes do conector no console do Realtime Compute for Apache Flink.

  3. Crie uma tabela Flink, construa dados Flink em tempo real usando um job Flink SQL e implante o job após o desenvolvimento.

    Na página de desenvolvimento de jobs Flink, crie e edite um job SQL. O exemplo abaixo define uma tabela de origem que gera dados aleatórios, uma tabela de resultados conectada ao MaxCompute e uma instrução INSERT para transferir os dados. Para mais informações sobre como desenvolver um job SQL, consulte Mapa de desenvolvimento de jobs.

    -- Create a Flink source table.
    CREATE TEMPORARY TABLE fake_src_table
    (
        c1 int,
        c2 VARCHAR,
        gt AS CURRENT_TIMESTAMP
    ) WITH (
      'connector' = 'faker',
      'fields.c2.expression' = '#{superhero.name}',
      'rows-per-second' = '100',
      'fields.c1.expression' = '#{number.numberBetween ''0'',''1000''}'
    );
    
    -- Create a temporary result table in Flink.
    CREATE TEMPORARY TABLE test_c_d_g 
    (
        c1 int,
        c2 VARCHAR,
        gt TIMESTAMP,
        ds varchar,
        PRIMARY KEY(c1) NOT ENFORCED
     ) PARTITIONED BY(ds)
     WITH (
        		'connector' = 'maxcompute',
        		'table.name' = 'mf_flink_upsert',
        		'sink.operation' = 'upsert',
        		'odps.access.id'='LTAI****************',
        		'odps.access.key'='********************',
        		'odps.end.point'='http://service.cn-beijing.maxcompute.aliyun.com/api',
        		'odps.project.name'='mf_mc_bj',
        		'upsert.write.bucket.num'='64'
    );
    
    -- Flink computation logic.
    INSERT INTO test_c_d_g
    SELECT  c1 AS c1,
            c2 AS c2,
            gt AS gt,
            date_format(gt, 'yyyyMMddHH') AS ds
    FROM    fake_src_table;

    Parâmetros:

    odps.end.point: Use o endpoint de rede interna da região correspondente.

    upsert.write.bucket.num: Este valor deve ser consistente com o valor da propriedade write.bucket.num da delta table criada no MaxCompute.

  4. Consulte a tabela MaxCompute para verificar se os dados foram gravados.

    SELECT * FROM mf_flink_upsert WHERE ds=2023061517;
    
    -- Your results may differ because the source data is randomly generated.
    +------+----+------+----+
    | c1   | c2 | gt   | ds |
    +------+----+------+----+
    | 0    | Skaar | 2023-06-16 01:59:41.116 | 2023061517 |
    | 21   | Supah Century | 2023-06-16 01:59:59.117 | 2023061517 |
    | 104  | Dark Gorilla Grodd | 2023-06-16 01:59:57.117 | 2023061517 |
    | 126  | Leader | 2023-06-16 01:59:39.116 | 2023061517 |
    

Apêndice: Parâmetros do conector Flink

  • Parâmetros básicos

    Parâmetro

    Obrigatório

    Valor padrão

    Descrição

    connector

    Sim

    Defina o tipo de conector como MaxCompute.

    odps.project.name

    Sim

    Nome do projeto MaxCompute.

    odps.access.id

    Sim

    AccessKey ID da sua conta RAM. Consulte a página AccessKey Pair.

    odps.access.key

    Sim

    AccessKey secret da sua conta RAM. Consulte a página AccessKey Pair.

    odps.end.point

    Sim

    Endpoint do MaxCompute. Para obter uma lista de endpoints regionais, consulte Endpoints.

    odps.tunnel.end.point

    Não

    Endpoint público do serviço Tunnel. Por padrão, as solicitações são roteadas automaticamente para o endpoint Tunnel apropriado. Defina este parâmetro para usar um endpoint específico e desativar o roteamento automático.

    Para mais informações sobre os endpoints do Tunnel em diferentes regiões e redes, consulte Endpoints.

    odps.tunnel.quota.name

    Não

    Nome da cota do Tunnel usada para acessar o MaxCompute.

    table.name

    Sim

    Nome da tabela MaxCompute no formato [project.][schema.]table.

    odps.namespace.schema

    Não

    false

    Define se o modelo de três camadas deve ser usado. Para mais informações sobre o modelo de três camadas, consulte Operações de schema.

    sink.operation

    Sim

    insert

    Tipo de gravação. Os valores válidos são insert ou upsert.

    Nota

    O modo upsert é compatível apenas com delta tables do MaxCompute.

    sink.parallelism

    Não

    Paralelismo de gravação. Se não definido, assume o paralelismo da fonte upstream como padrão.

    Nota

    Garanta que a propriedade da tabela write.bucket.num seja um múltiplo inteiro do valor configurado para otimizar o desempenho de gravação e maximizar a economia de memória no nó Sink.

    sink.meta.cache.time

    Não

    400

    Tamanho do cache de metadados.

    sink.meta.cache.expire.time

    Não

    1200

    Tempo de expiração do cache para metadados, em segundos.

    sink.coordinator.enable

    Não

    true

    Define se o modo coordenador deve ser ativado.

  • Parâmetros de partição

    Parâmetro

    Obrigatório

    Valor padrão

    Descrição

    sink.partition

    Não

    Nome da partição onde os dados serão gravados.

    Se você usar particionamento dinâmico, este parâmetro especifica o nome da partição pai das partições dinâmicas.

    sink.partition.default-value

    Não

    __DEFAULT_PARTITION__

    Nome padrão da partição quando o particionamento dinâmico é utilizado.

    sink.dynamic-partition.limit

    Não

    100

    Número máximo de partições que podem ser gravadas simultaneamente em um único checkpoint durante o particionamento dinâmico.

    Nota

    Aumentar significativamente este valor pode causar erros de falta de memória (OOM) no nó sink. Se o número de partições simultâneas exceder esse limite, o job falhará.

    sink.group-partition.enable

    Não

    false

    Define se o agrupamento por partição deve ser usado ao usar particionamento dinâmico.

    sink.partition.assigner.class

    Não

    Classe de implementação PartitionAssigner.

  • Parâmetros de gravação no modo FileCached

    Use o modo de cache de arquivo para jobs com um grande número de partições dinâmicas. Os parâmetros a seguir configuram esse modo.

    Parâmetro

    Obrigatório

    Valor padrão

    Descrição

    sink.file-cached.enable

    Não

    false

    Ativa o modo FileCached. Recomendado para jobs com muitas partições dinâmicas.

    • false: Modo FileCached desativado.

    • true: Modo FileCached ativado.

      Nota

      Quando o número de partições dinâmicas for alto, use o modo de cache de arquivo.

    sink.file-cached.tmp.dirs

    Não

    ./local

    Diretório padrão de cache de arquivos no modo FileCached.

    sink.file-cached.writer.num

    Não

    16

    Número de threads simultâneas de upload de dados para uma única tarefa no modo FileCached.

    Nota

    Não aumente significativamente o valor deste parâmetro. Se um número excessivo de partições for gravado simultaneamente, erros OOM provavelmente ocorrerão.

    sink.bucket.check-interval

    Não

    60000

    Intervalo para verificação do tamanho do arquivo no modo FileCached. Unidade: milissegundos (ms).

    sink.file-cached.rolling.max-size

    Não

    16 M

    Tamanho máximo para um único arquivo de cache.

    Quando um arquivo excede esse tamanho, ele é enviado.

    sink.file-cached.memory

    Não

    64 M

    Tamanho máximo de memória off-heap usado para gravação de arquivos no modo FileCached.

    sink.file-cached.memory.segment-size

    Não

    128 KB

    Tamanho do buffer usado para gravação de arquivos no modo FileCached.

    sink.file-cached.flush.always

    Não

    true

    Define se o cache deve ser usado para gravação de arquivos no modo FileCached.

    sink.file-cached.write.max-retries

    Não

    3

    Número de tentativas para upload de dados no modo FileCached.

  • Parâmetros de gravação Insert ou Upsert

    Parâmetros de gravação Upsert

    Parâmetro

    Obrigatório

    Valor padrão

    Descrição

    upsert.writer.max-retries

    Não

    3

    Número de tentativas após falha de um writer upsert ao gravar dados em um bucket.

    upsert.writer.buffer-size

    Não

    64 MB

    Tamanho do cache para um único writer upsert no Flink.

    Nota
    • Quando a soma dos tamanhos de buffer de todos os buckets atinge o limiar predefinido, o sistema aciona automaticamente uma operação de flush para atualizar os dados no servidor.

    • Um writer upsert grava dados em vários buckets simultaneamente. Recomendamos aumentar o valor deste parâmetro para melhorar a eficiência da gravação.

    • Se os dados forem gravados em um grande número de partições, podem ocorrer erros OOM. Nesse caso, reduza o valor deste parâmetro.

    upsert.writer.bucket.buffer-size

    Não

    1 MB

    Tamanho do cache para um único bucket no Flink. Se os recursos de memória no servidor Flink forem insuficientes, reduza o valor deste parâmetro.

    upsert.write.bucket.num

    Sim

    O número de buckets da tabela de destino deve ser igual ao valor de write.bucket.num.

    upsert.write.slot-num

    Não

    1

    Número de slots do Tunnel usados por uma única sessão.

    upsert.commit.max-retries

    Não

    3

    Número de tentativas para commit de uma sessão upsert.

    upsert.commit.thread-num

    Não

    16

    Paralelismo do commit de uma sessão upsert.

    Não defina este valor muito alto. Um número elevado de commits simultâneos aumenta o consumo de recursos, o que pode causar problemas de desempenho ou uso excessivo de recursos.

    upsert.major-compact.min-commits

    Não

    100

    Número mínimo de commits necessários para acionar uma compactação principal (major compaction).

    upsert.commit.timeout

    Não

    600

    Tempo limite para o commit de uma sessão upsert. Unidade: segundos (s).

    upsert.major-compact.enable

    Não

    false

    Define se a compactação principal deve ser ativada.

    upsert.flush.concurrent

    Não

    2

    Número máximo de buckets nos quais os dados podem ser gravados simultaneamente em uma única partição.

    Nota

    Quando os dados em um bucket passam por flush, um slot do Tunnel é ocupado.

    Nota

    Para mais informações sobre configurações recomendadas de parâmetros para gravações upsert, consulte Configurações recomendadas de parâmetros para gravações upsert.

    Parâmetros de gravação Insert

    Parâmetro

    Obrigatório

    Valor padrão

    Descrição

    insert.commit.thread-num

    Não

    16

    Paralelismo de uma sessão de commit.

    insert.arrow-writer.enable

    Não

    false

    Define se o formato Arrow deve ser usado.

    insert.arrow-writer.batch-size

    Não

    512

    Número máximo de linhas em um lote Arrow.

    insert.arrow-writer.flush-interval

    Não

    100000

    Intervalo de flush do writer. Unidade: milissegundos (ms).

    insert.writer.buffer-size

    Não

    64 MB

    Tamanho do cache do writer com buffer.