Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Object Storage Service

Última atualização: Jun 27, 2026

O Realtime Compute for Apache Flink permite ler e gravar no Object Storage Service (OSS) por meio do conector de sistema de arquivos. O OSS oferece durabilidade de dados de 99,9999999999% (doze noves) e disponibilidade de 99,995%, tornando-se um armazenamento confiável para pipelines Flink em grande escala.

Categoria

Detalhes

Tipos de tabela compatíveis

Tabelas de origem e de destino

Modos de execução

Batch e stream

Formatos de dados

ORC, Parquet, Avro, CSV, JSON e raw

Métricas de monitoramento específicas

Nenhuma

Tipos de API

DataStream API e SQL

Atualização ou exclusão de dados em tabelas de destino

Apenas inserção. Não há suporte para atualização nem exclusão.

Limitações

Gerais:

  • Somente o Ververica Runtime (VVR) 11 ou superior permite ler arquivos compactados (GZIP, BZIP2, XZ, DEFLATE) do OSS. O VVR 8 não processa corretamente arquivos compactados.

  • Versões do VVR anteriores à 8.0.6 aceitam apenas buckets do OSS na mesma conta. Para acessar buckets entre contas diferentes, use o VVR 8.0.6 ou posterior e configure a autenticação de bucket. Para mais detalhes, consulte Configurar autenticação de bucket.

  • Não há suporte para leitura incremental de novas partições.

  • Não há suporte para acesso ao OSS entre regiões. O bucket do OSS deve estar na mesma região do workspace do Flink. Como o conector do OSS se baseia na interface Filesystem, que aceita apenas um endpoint global único, tentar acessar um bucket em outra região gera um erro de incompatibilidade de endpoint.

Apenas para tabelas de destino:

Não há suporte para formatos de linha — Avro, CSV, JSON e raw — na gravação para o OSS. Consulte FLINK-30635 para obter mais detalhes.

Sintaxe

CREATE TABLE OssTable (
  column_name1 INT,
  column_name2 STRING,
  ...
  datetime STRING,
  `hour` STRING
) PARTITIONED BY (datetime, `hour`) WITH (
  'connector' = 'filesystem',           -- required: must be 'filesystem'
  'path' = 'oss://<bucket>/path',       -- required: URI of the OSS path
  'format' = '...',                     -- required: orc, parquet, avro, csv, json, or raw
  'partition.default-name' = '...',     -- optional: partition name when partition field is NULL or empty
  'source.monitor-interval' = '...',    -- optional (source only): interval to scan for new files
  'auto-compaction' = '...'            -- optional (sink only): enable automatic compaction after each checkpoint
);

Colunas de metadados

Tabelas de origem aceitam colunas de metadados que expõem informações no nível de arquivo para cada linha. Defina uma coluna de metadados na DDL adicionando METADATA após o tipo de dado:

CREATE TABLE MyUserTableWithFilepath (
  column_name1 INT,
  column_name2 STRING,
  `file.path` STRING NOT NULL METADATA
) WITH (
  'connector' = 'filesystem',
  'path' = 'oss://<bucket>/path',
  'format' = 'json'
)

As seguintes colunas de metadados estão disponíveis:

Chave

Tipo de dado

Descrição

file.path

STRING NOT NULL

Caminho completo do arquivo que contém a linha.

file.name

STRING NOT NULL

Nome do arquivo (último elemento do caminho).

file.size

BIGINT NOT NULL

Tamanho do arquivo, em bytes.

file.modification-time

TIMESTAMP_LTZ(3) NOT NULL

Última hora de modificação do arquivo.

Parâmetros WITH

Parâmetros gerais

Parâmetro

Obrigatório

Padrão

Descrição

connector

Sim

Deve ser filesystem.

path

Sim

Caminho do OSS no formato URI, como oss://my_bucket/my_path. No VVR 8.0.6 e versões posteriores, configure a autenticação de bucket após definir este parâmetro. Consulte Configurar autenticação de bucket.

format

Sim

Formato de arquivo: csv, json, avro, parquet, orc ou raw.

Parâmetros da tabela de origem

Parâmetro

Obrigatório

Padrão

Descrição

source.monitor-interval

Não

Intervalo para verificar novos arquivos. Deve ser maior que 0. Se não definido, o caminho é verificado apenas uma vez e a origem torna-se delimitada (bounded). Cada arquivo é identificado pelo caminho e processado exatamente uma vez. Os caminhos já processados são armazenados em estado e persistidos entre checkpoints e savepoints. Um intervalo menor acelera a descoberta de arquivos, mas aumenta a frequência de varredura.

Parâmetros da tabela de destino

Parâmetro

Obrigatório

Padrão

Descrição

partition.default-name

Não

_DEFAULT_PARTITION__

Nome da partição usado quando um campo de partição é NULL ou uma string vazia.

sink.rolling-policy.file-size

Não

128 MB

Tamanho máximo do arquivo parcial antes do rollover. Cada subtarefa de destino cria pelo menos um arquivo parcial por partição. Consulte Comportamento da política de rollover para entender a interação com os formatos de arquivo.

sink.rolling-policy.rollover-interval

Não

30min

Tempo máximo que um arquivo parcial pode permanecer aberto antes do rollover. A frequência de verificação é controlada por sink.rolling-policy.check-interval.

sink.rolling-policy.check-interval

Não

1min

Frequência de verificação para determinar se um arquivo parcial deve sofrer rollover com base em sink.rolling-policy.rollover-interval.

auto-compaction

Não

false

Define se a compactação automática deve ser ativada. Os dados são gravados primeiro em arquivos temporários. Após cada checkpoint, os arquivos temporários correspondentes são mesclados em arquivos maiores. Esses arquivos temporários não ficam visíveis antes da mesclagem. Quando ativado: apenas arquivos dentro de um mesmo checkpoint são mesclados (pelo menos um arquivo por checkpoint); a latência de visibilidade dos dados equivale a checkpoint interval + compaction duration; execuções longas de compactação podem causar backpressure e atrasar checkpoints.

compaction.file-size

Não

128 MB

Tamanho alvo do arquivo de saída compactada. Por padrão, assume o mesmo valor de sink.rolling-policy.file-size.

sink.partition-commit.trigger

Não

process-time

Define quando confirmar (commit) uma partição. Consulte Gatilhos de confirmação de partição.

sink.partition-commit.delay

Não

0s

Atraso mínimo antes de confirmar uma partição. Defina como 1 d para partições diárias ou 1 h para partições horárias.

sink.partition-commit.watermark-time-zone

Não

UTC

Fuso horário usado para converter uma watermark LONG em TIMESTAMP durante a comparação de confirmação de partição. Aplica-se apenas quando sink.partition-commit.trigger é partition-time. Use o fuso horário da sessão quando a watermark for definida em uma coluna TIMESTAMP_LTZ (por exemplo, Asia/Shanghai). Aceita nomes completos de fuso horário (como America/Los_Angeles) ou offsets personalizados (como GMT-08:00). Uma configuração incorreta pode atrasar a confirmação das partições em várias horas.

partition.time-extractor.kind

Não

default

Método de extração de tempo dos campos de partição. default: configura um padrão ou formatador de timestamp. custom: especifique uma classe extratora.

partition.time-extractor.class

Não

Classe que implementa a interface PartitionTimeExtractor. Obrigatória quando partition.time-extractor.kind é custom.

partition.time-extractor.timestamp-pattern

Não

Padrão para construir um timestamp a partir dos campos de partição. Por padrão, o primeiro campo é extraído usando yyyy-MM-dd hh:mm:ss. Exemplos: $dt (campo único), $year-$month-$day $hour:00:00 (vários campos), $dt $hour:00:00 (dois campos).

partition.time-extractor.timestamp-formatter

Não

yyyy-MM-dd HH:mm:ss

Formatador que converte a string de timestamp (definida por partition.time-extractor.timestamp-pattern) em um timestamp. Por exemplo, se partition.time-extractor.timestamp-pattern for $year$month$day, defina este parâmetro como yyyyMMdd. Compatível com o DateTimeFormatter do Java.

sink.partition-commit.policy.kind

Não

Forma de notificar consumidores downstream de que uma partição está pronta. success-file: grava um arquivo _SUCCESS no diretório da partição. custom: utiliza uma classe que implementa PartitionCommitPolicy. É possível combinar várias políticas.

sink.partition-commit.policy.class

Não

Classe que implementa PartitionCommitPolicy. Obrigatória quando sink.partition-commit.policy.kind é custom.

sink.partition-commit.success-file.name

Não

_SUCCESS

Nome do arquivo de sucesso gravado pela política de confirmação success-file.

sink.parallelism

Não

Paralelismo do operador de gravação de arquivos. Assume por padrão o paralelismo do operador upstream. Deve ser maior que 0. Quando auto-compaction está ativado, o operador de compactação também utiliza esse paralelismo.

Comportamento da política de rollover

O comportamento de rollover varia conforme o formato do arquivo:

Em formatos colunares (Parquet, ORC, Avro), o arquivo parcial sempre sofre rollover no momento do checkpoint, mesmo que os critérios da política de rollover não tenham sido atendidos. O tamanho do arquivo e o intervalo de rollover funcionam como gatilhos adicionais entre checkpoints.
Em formatos de linha (CSV, JSON, raw), o arquivo parcial só sofre rollover quando os critérios da política são atendidos ( sink.rolling-policy.file-size ou sink.rolling-policy.rollover-interval ). Para garantir baixa latência na visibilidade dos arquivos, ajuste sink.rolling-policy.rollover-interval em conjunto com o intervalo de checkpoint.
Não há suporte para formatos de linha em tabelas de destino do OSS devido ao problema FLINK-30635 . O comportamento descrito acima será aplicável caso o suporte a formatos de linha seja adicionado em uma versão futura.

Gatilhos de confirmação de partição

Dois tipos de gatilho estão disponíveis para sink.partition-commit.trigger:

  • process-time (padrão): Confirma a partição quando a hora atual do sistema ultrapassa o horário de criação da partição somado ao valor de sink.partition-commit.delay. Não exige watermark nem extrator de tempo de partição. É mais genérico, porém menos preciso — atrasos ou falhas nos dados podem provocar confirmações prematuras.

  • partition-time: Confirma a partição quando a watermark ultrapassa o horário de criação da partição somado ao valor de sink.partition-commit.delay. Requer geração de watermark e partições baseadas em tempo (horárias, diárias, etc.).

Configurar autenticação de bucket

Somente o VVR 8.0.6 e versões posteriores oferecem suporte à autenticação de bucket.

Após definir o parâmetro path, configure a autenticação de bucket para que o Flink possa ler e gravar no caminho especificado do OSS. Adicione as configurações abaixo na seção Additional Configurations, na aba Parameters da página Deployment Details no console de desenvolvimento do Realtime Compute:

fs.oss.bucket.<bucketName>.accessKeyId: <your-access-key-id>
fs.oss.bucket.<bucketName>.accessKeySecret: <your-access-key-secret>

Substitua <bucketName> pelo nome do bucket informado no parâmetro path.

Item de configuração

Descrição

fs.oss.bucket.<bucketName>.accessKeyId

AccessKey ID do bucket. Use um AccessKey existente ou crie um novo. Consulte Criar um AccessKey.

fs.oss.bucket.<bucketName>.accessKeySecret

AccessKey Secret do bucket.

Importante

O AccessKey Secret é exibido apenas uma vez no momento da criação. Armazene-o com segurança.

Gravar no OSS-HDFS

Adicione a seguinte configuração na seção Additional Configurations, na aba Parameters da página Deployment Details no console de desenvolvimento do Realtime Compute:

fs.oss.jindo.buckets: <bucket-names>
fs.oss.jindo.accessKeyId: <your-access-key-id>
fs.oss.jindo.accessKeySecret: <your-access-key-secret>

Item de configuração

Descrição

fs.oss.jindo.buckets

Nomes dos buckets OSS-HDFS, separados por ponto e vírgula. Quando o Flink grava em um caminho do OSS e o bucket correspondente estiver listado aqui, os dados serão direcionados ao serviço OSS-HDFS.

fs.oss.jindo.accessKeyId

AccessKey ID. Consulte Criar um AccessKey.

fs.oss.jindo.accessKeySecret

AccessKey Secret.

Importante

O AccessKey Secret é exibido apenas uma vez no momento da criação. Armazene-o com segurança.

Configure o endpoint do OSS-HDFS usando um dos métodos a seguir:

Parameter configuration

Adicione o endpoint em Additional Configurations:

fs.oss.jindo.endpoint: <oss-hdfs-endpoint>

Path configuration

Incorpore o endpoint diretamente no caminho do OSS:

oss://<bucket-name>.<oss-hdfs-endpoint>/<directory>

Ao usar este método, fs.oss.jindo.buckets deve incluir <bucket-name>.<oss-hdfs-endpoint>.

Por exemplo, se o nome do bucket for jindo-test e o endpoint for cn-beijing.oss-dls.aliyuncs.com:

# OSS path
oss://jindo-test.cn-beijing.oss-dls.aliyuncs.com/<directory>

# Additional Configurations
fs.oss.jindo.buckets: jindo-test,jindo-test.cn-beijing.oss-dls.aliyuncs.com

Gravação em um Hadoop Distributed File System (HDFS) externo

Para caminhos que usam o esquema hdfs://, adicione a configuração abaixo para especificar ou alterar o nome de usuário de acesso:

containerized.taskmanager.env.HADOOP_USER_NAME: hdfs
containerized.master.env.HADOOP_USER_NAME: hdfs

Exemplos

Leitura do OSS (tabela de origem)

CREATE TEMPORARY TABLE fs_table_source (
  `id` INT,
  `name` VARCHAR
) WITH (
  'connector' = 'filesystem',
  'path' = 'oss://<bucket>/path',
  'format' = 'parquet'
);

CREATE TEMPORARY TABLE blackhole_sink (
  `id` INT,
  `name` VARCHAR
) WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink SELECT * FROM fs_table_source;

Gravação no OSS (tabela de destino)

Gravar em uma tabela particionada

Este exemplo transmite dados de uma origem datagen, particiona-os por data e hora e confirma as partições usando o gatilho partition-time:

CREATE TABLE datagen_source (
  user_id STRING,
  order_amount DOUBLE,
  ts BIGINT,                                               -- Timestamp in milliseconds
  ts_ltz AS TO_TIMESTAMP_LTZ(ts, 3),
  WATERMARK FOR ts_ltz AS ts_ltz - INTERVAL '5' SECOND    -- Watermark on TIMESTAMP_LTZ column
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE fs_table_sink (
  user_id STRING,
  order_amount DOUBLE,
  dt STRING,
  `hour` STRING
) PARTITIONED BY (dt, `hour`) WITH (
  'connector' = 'filesystem',
  'path' = 'oss://<bucket>/path',
  'format' = 'parquet',
  'partition.time-extractor.timestamp-pattern' = '$dt $hour:00:00',
  'sink.partition-commit.delay' = '1 h',
  'sink.partition-commit.trigger' = 'partition-time',
  'sink.partition-commit.watermark-time-zone' = 'Asia/Shanghai',
  'sink.partition-commit.policy.kind' = 'success-file'
);

INSERT INTO fs_table_sink
SELECT
  user_id,
  order_amount,
  DATE_FORMAT(ts_ltz, 'yyyy-MM-dd'),
  DATE_FORMAT(ts_ltz, 'HH')
FROM datagen_source;

Gravar em uma tabela não particionada

CREATE TABLE datagen_source (
  user_id STRING,
  order_amount DOUBLE
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE fs_table_sink (
  user_id STRING,
  order_amount DOUBLE
) WITH (
  'connector' = 'filesystem',
  'path' = 'oss://<bucket>/path',
  'format' = 'parquet'
);

INSERT INTO fs_table_sink SELECT * FROM datagen_source;

DataStream API

Importante

Para usar a DataStream API, configure primeiro o conector DataStream. Consulte Usar um conector DataStream.

O exemplo a seguir usa StreamingFileSink com OnCheckpointRollingPolicy para gravar no OSS. Os arquivos parciais sofrem rollover a cada checkpoint.

String outputPath = "oss://<bucket>/path";

final StreamingFileSink<Row> sink =
    StreamingFileSink.forRowFormat(
            new Path(outputPath),
            (Encoder<Row>) (element, stream) -> {
                out.println(element.toString());
            })
        .withRollingPolicy(OnCheckpointRollingPolicy.build())
        .build();

outputStream.addSink(sink);

Para gravar no OSS-HDFS, configure também os parâmetros correspondentes em Additional Configurations. Consulte Gravar no OSS-HDFS.

Próximos passos