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 |
|
|
STRING NOT NULL |
Caminho completo do arquivo que contém a linha. |
|
|
STRING NOT NULL |
Nome do arquivo (último elemento do caminho). |
|
|
BIGINT NOT NULL |
Tamanho do arquivo, em bytes. |
|
|
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 |
|
|
Sim |
— |
Deve ser |
|
|
Sim |
— |
Caminho do OSS no formato URI, como |
|
|
Sim |
— |
Formato de arquivo: |
Parâmetros da tabela de origem
|
Parâmetro |
Obrigatório |
Padrão |
Descrição |
|
|
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 |
|
|
Não |
|
Nome da partição usado quando um campo de partição é NULL ou uma string vazia. |
|
|
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. |
|
|
Não |
30min |
Tempo máximo que um arquivo parcial pode permanecer aberto antes do rollover. A frequência de verificação é controlada por |
|
|
Não |
1min |
Frequência de verificação para determinar se um arquivo parcial deve sofrer rollover com base em |
|
|
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 |
|
|
Não |
128 MB |
Tamanho alvo do arquivo de saída compactada. Por padrão, assume o mesmo valor de |
|
|
Não |
|
Define quando confirmar (commit) uma partição. Consulte Gatilhos de confirmação de partição. |
|
|
Não |
|
Atraso mínimo antes de confirmar uma partição. Defina como |
|
|
Não |
|
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 |
|
|
Não |
|
Método de extração de tempo dos campos de partição. |
|
|
Não |
— |
Classe que implementa a interface |
|
|
Não |
— |
Padrão para construir um timestamp a partir dos campos de partição. Por padrão, o primeiro campo é extraído usando |
|
|
Não |
|
Formatador que converte a string de timestamp (definida por |
|
|
Não |
— |
Forma de notificar consumidores downstream de que uma partição está pronta. |
|
|
Não |
— |
Classe que implementa |
|
|
Não |
|
Nome do arquivo de sucesso gravado pela política de confirmação |
|
|
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 |
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-sizeousink.rolling-policy.rollover-interval). Para garantir baixa latência na visibilidade dos arquivos, ajustesink.rolling-policy.rollover-intervalem 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 desink.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 desink.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 |
|
|
AccessKey ID do bucket. Use um AccessKey existente ou crie um novo. Consulte Criar um AccessKey. |
|
|
AccessKey Secret do bucket. |
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 |
|
|
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. |
|
|
AccessKey ID. Consulte Criar um AccessKey. |
|
|
AccessKey Secret. |
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
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.