Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:9 de abril de 2026

Última atualização: Jun 27, 2026

O VVR 11.6.0 adiciona funções de IA multimodal (PDFs e imagens), o tipo Variant para dados semiestruturados e promove a Ingestão de Dados (Flink CDC) para disponibilidade geral (GA). Esta versão também aprimora os conectores Kafka, MySQL CDC, OceanBase, Elasticsearch e Hologres, além de incorporar correções das comunidades Apache Flink 1.20.2 e 1.20.3.

Importante

Esta atualização está sendo implementada em fases em todas as regiões. Se sua conta ainda não foi atualizada, esses recursos estarão indisponíveis. Para solicitar uma atualização antecipada, envie um ticket.

Aprimoramentos do mecanismo

IA

Multimodal

Novas funções integradas para processamento de imagens e documentos em tempo real:

  • Conversão de PDF para imagem: converta páginas de PDF em imagens dentro de um job do Flink.

  • Recuperação de arquivos: leia o conteúdo de arquivos diretamente do OSS e do Message Notification Service (MNS).

  • Detecção de nitidez de imagem: avalie a nitidez da imagem com OpenCV.

  • Compressão de imagem: reduza o tamanho da imagem inline.

  • Repasse de imagem Base64: encaminhe imagens codificadas em Base64 para chamadas de modelos downstream.

Essas funções oferecem suporte a pipelines de inferência ponta a ponta com Vision Language Models (VLMs), como o Qwen-VL.

Conector Simple Message Queue (MNS)

Um novo conector MNS permite assinar eventos de alteração do OSS e fechar o ciclo em pipelines multimodais em tempo real.

Flink SQL

Tipo Variant

O novo tipo VARIANT lida com dados semiestruturados sem esquema fixo. Acesse campos com notação de ponto (variant.field) ou colchetes (variant['key']), converta tipos básicos e grave colunas Variant em sinks Paimon.

Novas funções

  • parse_json: converte uma string JSON em um valor Variant (disponível no Transform do Flink CDC).

  • MD5 e outras funções de hash: disponíveis no Transform do Flink CDC.

Ingestão de Dados (Flink CDC) — agora GA

A Ingestão de Dados (Flink CDC) saiu da prévia pública e agora tem disponibilidade geral (GA).

Status de GA dos conectores

Status

Conectores

Disponibilidade geral

Paimon, StarRocks, Hologres, MySQL, Kafka

Prévia pública

Doris, OceanBase, MaxCompute, SLS, MongoDB, Postgres

Novidades nesta versão

Mesclagem de colunas

Combine vários campos upstream com nomes ou capitalizações diferentes em uma única coluna de destino. Esse recurso oferece suporte a correspondência por regex, normalização de maiúsculas/minúsculas e mapeamentos personalizados, sendo útil para consolidar fontes JSON com nomenclatura de campos inconsistente.

Gravações apenas de acréscimo em tabelas particionadas do Paimon

O sink Paimon agora grava em tabelas particionadas sem chave primária. Anteriormente, a chave de partição precisava fazer parte da chave primária; esse requisito foi removido para casos de uso apenas de acréscimo.

Se você depende atualmente do comportamento antigo, verifique a configuração do seu sink Paimon após a atualização.

Aprimoramentos do Transform

  • Limpe uma chave primária ou de partição passando um valor null.

  • Defina roteamentos complexos de nomes de tabela usando expressões regulares.

Suporte a Variant

  • Acesse e converta campos do tipo VARIANT; grave dados Variant no Paimon.

Aprimoramentos de source

  • Source PolarDB-X CDC (prévia pública): assinatura de binlog de alto paralelismo no nível da tabela e assinatura por dimensão de tabela.

  • Source SLS: imponha tipos específicos de análise de campos.

  • Source Kafka: divida uma única mensagem do Kafka em vários registros com base nos valores dos campos (roteamento de campo) e grave-os em diferentes tabelas de destino. Particionadores personalizados também são suportados.

Aprimoramentos de sink

Sink

Aprimoramento

Paimon

Configuração separada de paralelismo para nós de commit

MaxCompute

Mapeamento de tipo DATETIME; lógica de commit otimizada para reduzir o consumo de consultas por segundo (QPS)

Iceberg

Referências de catálogo integradas; recuperação automática de URLs de conexão e credenciais

Conectores

Kafka

  • O sink grava IDs de tabela de três partes (Database.Schema.Table) no formato Debezium JSON.

  • Após uma alteração de tópico, uma reinicialização com estado agora gera uma exceção de estado incompatível em vez de consumir silenciosamente dos tópicos antigo e novo simultaneamente. Isso evita consumo duplo e inconsistências de dados.

Se seus jobs dependem de reinicialização após uma mudança de tópico, revise esse comportamento antes de atualizar.

MySQL CDC

  • As mensagens de erro para identificadores globais de transação (GTIDs) expirados agora indicam claramente a causa raiz.

  • O ID do servidor consumidor é incluído nos logs para simplificar a solução de problemas.

PolarDB-X

Agora há suporte como source YAML CDC (prévia pública).

OceanBase

O sink JDBC agora oferece suporte a rollbacks manuais de transação e reutilização de pool de conexões, resolvendo desconexões frequentes causadas por wait_timeout.

Elasticsearch

A tabela source e a tabela de dimensão agora oferecem suporte ao Elasticsearch 8.x (com o cliente compatível com ES7).

Doris

As mensagens de erro para configurações incorretas de porta estão mais descritivas.

Integração com data lakehouse

Iceberg

  • O sink agora oferece suporte à métrica numRecordsOutOfSinkPerSecond (OUT RPS).

  • Configure parâmetros relacionados ao Hadoop para maior flexibilidade de conexão.

  • Jobs do Flink CDC podem gravar no Iceberg do Data Lake Formation (DLF).

Hologres

  • A tabela source Binlog permite consumo a partir do offset LATEST.

  • O catálogo Hologres expõe índices secundários e chaves de varredura por prefixo.

  • A leitura do tipo de array varchar[] agora é suportada.

  • O cache de detecção de parâmetros foi otimizado para evitar tempos limite de inicialização quando há muitas tabelas.

  • O sink aceita paralelismo superior ao número de shards quando sink.reshuffle-by-holo-distribution-key.enabled está configurado.

MaxCompute

  • O catálogo usa consultas paginadas para evitar o congelamento do centro de metadados.

  • A lógica de commit do sink YAML otimizada reduz erros OOM causados por limites de QPS.

Hive

O catálogo permite especificar um formato de armazenamento (como Parquet) ao criar uma tabela.

Paimon

Adicionado suporte ao formato de arquivo Lance.

Observabilidade

Novas métricas:

Métrica

Descrição

geminiDB.disk_space_*

Uso de disco local

geminiDB.native_memory_usage

Memória nativa do GeminiDB em uso

geminiDB.native_memory_limit

Limite de memória nativa do GeminiDB

sourceParallelismUpperBound

Limite superior de paralelismo do operador auto-pilot

Logs WARN não essenciais — como alertas quando um formato não oferece suporte a snapshots — agora são suprimidos.

Correções de bugs

Estabilidade

  • Correções da comunidade Apache Flink: incorpora correções das versões 1.20.2 e 1.20.3 do Apache Flink.

  • Perda de dados no Kafka: corrigida a perda de dados ao ler do Kafka e gravar no OSS com transações ativadas.

  • Desconexão do PolarDB-X: corrigido um pico abrupto de latência e a exceção EOFException acionados por uma conexão interrompida do PolarDB-X.

  • wait_timeout do OceanBase: corrigidas desconexões frequentes no sink JDBC do OceanBase causadas por wait_timeout.

Corretude

  • Canal Protobuf do Flink CDC: corrigido o formato de timestamp inconsistente e o tratamento do tipo tinyint.

  • Depuração do source MySQL CDC: corrigida uma contagem com erro de um dígito quando a reutilização estava ativada (a contagem de tabelas agora exibe o número correto).

  • Sink MaxCompute (ODPS) do Flink CDC: corrigidos erros OOM de Metaspace causados por commits frequentes.

Experiência

  • Mensagens de erro de Temporal Join: as mensagens de erro para sintaxe de Temporal Join estão mais claras.

  • Logs WARN internos: Cannot snapshot the table e alertas internos semelhantes agora são registrados no nível DEBUG.

  • Campos nulos no binlog do Hologres: corrigida uma exceção que ocorria quando alguns campos eram nulos durante o consumo de binlog do Hologres.