A API DataStream do Flink permite definir transformações de dados personalizadas, operadores e lógica de processamento para workloads de streaming complexos. Este tópico orienta o desenvolvimento de um job JAR com a API DataStream, desde a configuração das dependências do Maven até o empacotamento e a implantação do job.
Pré-requisitos
Antes de começar, verifique se você tem:
Um ambiente de desenvolvimento integrado (IDE), como o IntelliJ IDEA
Maven 3.6.3 ou posterior
JDK 8 ou JDK 11 (nenhuma outra versão é suportada)
As fontes de dados necessárias para o job. Este exemplo usa o Message Queue for Apache Kafka 2.6.2 e o ApsaraDB RDS for MySQL 8.0.
As instâncias do Kafka e do MySQL devem estar na mesma Virtual Private Cloud (VPC) do workspace do Realtime Compute for Apache Flink. As instâncias do MySQL também precisam estar na mesma região. Caso não tenha essas fontes de dados, configure-as antes de prosseguir: Criar uma instância do Message Queue for Apache Kafka e Criar uma instância do ApsaraDB RDS for MySQL . Para acesso entre VPCs ou pela internet, consulte Selecionar um método de conexão de rede .
Os jobs JAR são desenvolvidos offline e implantados no console do Realtime Compute for Apache Flink. A API DataStream no Realtime Compute for Apache Flink é totalmente compatível com o Apache Flink open source. Para mais contexto, consulte a visão geral da arquitetura do Apache Flink e o guia do desenvolvedor da API DataStream do Flink.
Etapa 1: Configurar as dependências do Flink
Adicione as dependências principais do Flink ao arquivo POM do Maven. Essas dependências são necessárias em tempo de compilação, mas não devem ser empacotadas no JAR do job, pois o ambiente de execução já as fornece. Incluí-las resulta, no melhor dos casos, em um JAR excessivamente grande e, no pior, em conflitos de carregamento de classes entre as versões empacotadas e as do ambiente de execução.
Defina o escopo como provided para todas as dependências org.apache.flink que começam com flink-:
<!-- Core Flink dependencies — set to provided to avoid JAR conflicts -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
O valor de ${flink.version} deve corresponder à versão do Flink do mecanismo Ververica Runtime (VVR) selecionado durante a implantação do job. Por exemplo, o mecanismo vvr-8.0.9-flink-1.17 corresponde ao Flink 1.17.2. Apenas métodos anotados com @Public ou @PublicEvolving no código-fonte do Flink são considerados estáveis. O Realtime Compute for Apache Flink garante compatibilidade somente para esses métodos. Para consultar detalhes sobre as versões do mecanismo VVR, visualize Como visualizo a versão do Flink do job atual?
Para obter a lista completa de dependências, incluindo logging, consulte o código de exemplo completo ao final deste tópico.
Etapa 2: Adicionar dependências de conectores
Para ler e gravar em sistemas externos com a API DataStream, adicione o conector DataStream do VVR correspondente. Use apenas conectores listados explicitamente como compatíveis com a API DataStream em Conectores suportados. Outros conectores podem alterar interfaces sem aviso prévio.
Escolha uma estratégia de importação
|
Estratégia |
Escopo da dependência |
Quando usar |
|
Carregar o conector como dependência adicional (recomendado) |
|
Mantém o JAR do job leve e permite gerenciar versões de conectores independentemente. O JAR do conector é carregado separadamente em tempo de execução. |
|
Empacotar o conector no JAR do job |
|
Reúne tudo em um único artefato. Recomendado quando é necessário controle rigoroso sobre dependências transitivas ou quando uma etapa separada de upload não é viável. |
Opção 1: Carregar como dependência adicional (recomendado)
-
Adicione o conector como uma dependência
providedno arquivo POM. Como o JAR do conector é carregado separadamente em tempo de execução, definirprovidedimpede que ele seja incluído no JAR do job.<!-- Kafka connector — loaded as an additional dependency at runtime --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-kafka</artifactId> <version>${vvr.version}</version> <scope>provided</scope> </dependency> <!-- MySQL connector — loaded as an additional dependency at runtime --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-mysql</artifactId> <version>${vvr.version}</version> <scope>provided</scope> </dependency>${vvr.version}representa a versão do mecanismo de execução VVR. Por exemplo, se o job executa emvvr-8.0.9-flink-1.17, a versão correspondente do Flink é1.17.2. Sempre que possível, utilize a versão mais recente do mecanismo. Consulte Mecanismos. -
Caso esteja desenvolvendo um conector personalizado ou estendendo um existente, adicione também estes pacotes base:
<!-- Base interface for Flink connectors --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-base</artifactId> <version>${flink.version}</version> </dependency> <!-- Base interface for Alibaba Cloud connectors --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-common</artifactId> <version>${vvr.version}</version> </dependency> Para configurações de conexão DataStream e exemplos de código, consulte a documentação de cada conector. A lista completa de conectores compatíveis com DataStream está disponível em Conectores suportados.
-
Ao implantar o job, faça o upload do arquivo JAR do conector na seção Additional Dependencies. Baixe os conectores fornecidos pela Alibaba Cloud na lista de conectores. Para etapas de implantação, consulte Implantar um job JAR.

Opção 2: Empacotar o conector no JAR do job
-
Adicione os conectores com o escopo padrão
compile. Eles serão incluídos diretamente no JAR do job.<!-- Kafka connector — bundled into the job JAR --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-kafka</artifactId> <version>${vvr.version}</version> </dependency> <!-- MySQL connector — bundled into the job JAR --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-mysql</artifactId> <version>${vvr.version}</version> </dependency> Se estiver desenvolvendo ou estendendo um conector, adicione os mesmos pacotes base descritos na Opção 1.
Para configuração de conexão DataStream, consulte a documentação de cada conector em Conectores suportados.
Etapa 3: Ler configurações do OSS
Jobs JAR não leem arquivos de configuração locais diretamente da função main. Em vez disso, carregue o arquivo de configuração em um bucket do Object Storage Service (OSS) associado ao workspace do Flink e leia-o em tempo de execução.
Armazenar credenciais em um arquivo de configuração, em vez de codificá-las diretamente no código-fonte, é uma prática recomendada de segurança.
-
Crie um arquivo
config.properties:# Kafka bootstrapServers=host1:9092,host2:9092,host3:9092 inputTopic=topic groupId=groupId # MySQL database.url=jdbc:mysql://localhost:3306/my_database database.username=username database.password=password Carregue este arquivo no código do job utilizando um dos métodos a seguir.
Método 1: Ler do bucket OSS do workspace
Faça o upload do arquivo pela página Resource Management no console do Realtime Compute for Apache Flink. Quando o job for executado, os arquivos de dependência adicionais serão carregados no diretório /flink/usrlib/ do pod do job.
Properties properties = new Properties();
Map<String, String> configMap = new HashMap<>();
try (InputStream input = new FileInputStream("/flink/usrlib/config.properties")) {
// Load the properties file.
properties.load(input);
// Read property values.
configMap.put("bootstrapServers", properties.getProperty("bootstrapServers"));
configMap.put("inputTopic", properties.getProperty("inputTopic"));
configMap.put("groupId", properties.getProperty("groupId"));
configMap.put("url", properties.getProperty("database.url"));
configMap.put("username", properties.getProperty("database.username"));
configMap.put("password", properties.getProperty("database.password"));
} catch (IOException ex) {
ex.printStackTrace();
}
Método 2: Ler diretamente de um bucket OSS
Carregue o arquivo em qualquer bucket OSS ao qual o workspace tenha permissão de acesso e, em seguida, leia-o usando OSSClient. Para mais informações, consulte Download via stream e Gerenciar credenciais de acesso.
OSS ossClient = new OSSClientBuilder().build("Endpoint", "AccessKeyId", "AccessKeySecret");
try (OSSObject ossObject = ossClient.getObject("examplebucket", "exampledir/config.properties");
BufferedReader reader = new BufferedReader(new InputStreamReader(ossObject.getObjectContent()))) {
// Read file and process...
} finally {
if (ossClient != null) {
ossClient.shutdown();
}
}
Etapa 4: Escrever a lógica de negócios
Os trechos a seguir mostram os dois principais padrões DataStream usados no exemplo completo: integração de uma fonte e transformação de dados.
Integrar uma fonte Kafka. Uma watermark mede o progresso do tempo de evento e é frequentemente usada com timestamps. Este exemplo não utiliza política de watermark.
// Attach a Kafka source to the execution environment.
// WatermarkStrategy.noWatermarks() skips event-time watermark processing.
DataStreamSource<String> stream = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "kafka Source");
Para opções de watermark, consulte Estratégias de Watermark.
Transformar o stream. Este exemplo converte DataStream<String> para DataStream<Student> e filtra registros com pontuação inferior a 60.
// Map each comma-separated string to a Student object, then filter by score.
DataStream<Student> source = stream
.map(new MapFunction<String, Student>() {
@Override
public Student map(String s) throws Exception {
String[] data = s.split(",");
return new Student(Integer.parseInt(data[0]), data[1], Integer.parseInt(data[2]));
}
}).filter(student -> student.score >= 60);
Para mais tipos de operadores, consulte Operadores do Flink.
Etapa 5: Empacotar o job
Use o maven-shade-plugin para criar um fat JAR que agrupe todas as dependências necessárias.
Regras de escopo durante o empacotamento:
Conectores carregados como dependências adicionais: mantenha o escopo como
provided— não os inclua no pacote.Conectores empacotados no JAR: use o escopo padrão
compile.
<build>
<plugins>
<!-- Java compiler plugin -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.11.0</version>
<configuration>
<source>${target.java.version}</source>
<target>${target.java.version}</target>
</configuration>
</plugin>
<!-- Shade plugin: creates a fat JAR with all required dependencies.
Update <mainClass> if your program entry point changes. -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.2.0</version>
<executions>
<execution>
<phase>package</phase>
<goals>
<goal>shade</goal>
</goals>
<configuration>
<artifactSet>
<excludes>
<!-- Exclude unnecessary dependencies -->
<exclude>org.apache.flink:force-shading</exclude>
<exclude>com.google.code.findbugs:jsr305</exclude>
<exclude>org.slf4j:*</exclude>
<exclude>org.apache.logging.log4j:*</exclude>
</excludes>
</artifactSet>
<filters>
<filter>
<!-- Exclude META-INF signatures to prevent security exceptions. -->
<artifact>*:*</artifact>
<excludes>
<exclude>META-INF/*.SF</exclude>
<exclude>META-INF/*.DSA</exclude>
<exclude>META-INF/*.RSA</exclude>
</excludes>
</filter>
</filters>
<transformers>
<transformer
implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
<mainClass>com.aliyun.FlinkDemo</mainClass>
</transformer>
</transformers>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
Etapa 6: Testar e implantar
O Realtime Compute for Apache Flink não tem acesso à internet por padrão. Isso impede testes locais ponta a ponta com conectores ativos na maioria das configurações. Execute testes unitários da lógica de negócios independentemente antes da implantação. Para orientações sobre como executar testes dependentes de conectores localmente, consulte Executar e depurar um job com conectores localmente.
Para implantar o job:
Siga as etapas em Implantar um job JAR.
-
Se estiver utilizando a abordagem de dependência adicional, faça o upload do arquivo JAR do conector e do arquivo
config.propertiesna seção Additional Dependencies.
Configure os parâmetros de execução — checkpoints, Time to Live (TTL) e estratégias de reinicialização — na página Deployment Details após a implantação. Configurações definidas na página facilitam atualizações futuras sem necessidade de modificar ou reempacotar o código do job. Observe que configurações no nível de código têm prioridade sobre as definições da página. Para detalhes, consulte Configurar informações de implantação do job.
Código de exemplo completo
Este exemplo lê registros de alunos do Kafka, filtra pontuações abaixo de 60 e grava os registros aprovados em uma tabela MySQL. Para diretrizes de estilo de código, consulte o guia de estilo e qualidade de código do Flink.
Parâmetros de execução como checkpoints, TTL e estratégias de reinicialização não estão incluídos neste exemplo. Configure-os na página Deployment Details após a implantação para simplificar atualizações futuras. Para detalhes, consulte Configurar informações de implantação do job .
FlinkDemo.java
package com.aliyun;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.connector.jdbc.JdbcConnectionOptions;
import org.apache.flink.connector.jdbc.JdbcExecutionOptions;
import org.apache.flink.connector.jdbc.JdbcSink;
import org.apache.flink.connector.jdbc.JdbcStatementBuilder;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;
import org.apache.flink.kafka.shaded.org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import java.io.FileInputStream;
import java.io.IOException;
import java.io.InputStream;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.util.HashMap;
import java.util.Map;
import java.util.Properties;
public class FlinkDemo {
// Data model for student records.
public static class Student {
public int id;
public String name;
public int score;
public Student(int id, String name, int score) {
this.id = id;
this.name = name;
this.score = score;
}
}
public static void main(String[] args) throws Exception {
// Create a Flink execution environment.
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
Properties properties = new Properties();
Map<String, String> configMap = new HashMap<>();
// Load connection configuration from the file uploaded as an additional dependency.
// The file is available at /flink/usrlib/ when the job runs.
try (InputStream input = new FileInputStream("/flink/usrlib/config.properties")) {
properties.load(input);
configMap.put("bootstrapServers", properties.getProperty("bootstrapServers"));
configMap.put("inputTopic", properties.getProperty("inputTopic"));
configMap.put("groupId", properties.getProperty("groupId"));
configMap.put("url", properties.getProperty("database.url"));
configMap.put("username", properties.getProperty("database.username"));
configMap.put("password", properties.getProperty("database.password"));
} catch (IOException ex) {
ex.printStackTrace();
}
// Build the Kafka source starting from the latest offset.
KafkaSource<String> kafkaSource = KafkaSource.<String>builder()
.setBootstrapServers(configMap.get("bootstrapServers"))
.setTopics(configMap.get("inputTopic"))
.setStartingOffsets(OffsetsInitializer.latest())
.setGroupId(configMap.get("groupId"))
.setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class))
.build();
// Attach the Kafka source. No watermark policy is applied in this example.
DataStreamSource<String> stream = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "kafka Source");
// Parse each record into a Student object and keep only scores >= 60.
DataStream<Student> source = stream
.map(new MapFunction<String, Student>() {
@Override
public Student map(String s) throws Exception {
String[] data = s.split(",");
return new Student(Integer.parseInt(data[0]), data[1], Integer.parseInt(data[2]));
}
}).filter(Student -> Student.score >= 60);
// Write passing records to MySQL in batches of 5, flushing every 2 seconds.
source.addSink(JdbcSink.sink(
"INSERT IGNORE INTO student (id, username, score) VALUES (?, ?, ?)",
new JdbcStatementBuilder<Student>() {
public void accept(PreparedStatement ps, Student data) {
try {
ps.setInt(1, data.id);
ps.setString(2, data.name);
ps.setInt(3, data.score);
} catch (SQLException e) {
throw new RuntimeException(e);
}
}
},
new JdbcExecutionOptions.Builder()
.withBatchSize(5) // Records per batch.
.withBatchIntervalMs(2000) // Max flush interval in milliseconds.
.build(),
new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
.withUrl(configMap.get("url"))
.withDriverName("com.mysql.cj.jdbc.Driver")
.withUsername(configMap.get("username"))
.withPassword(configMap.get("password"))
.build()
)).name("Sink MySQL");
env.execute("Flink Demo");
}
}
pom.xml
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/maven-v4_0_0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.aliyun</groupId>
<artifactId>FlinkDemo</artifactId>
<version>1.0-SNAPSHOT</version>
<name>FlinkDemo</name>
<packaging>jar</packaging>
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<flink.version>1.17.1</flink.version>
<vvr.version>1.17-vvr-8.0.4-1</vvr.version>
<target.java.version>1.8</target.java.version>
<maven.compiler.source>${target.java.version}</maven.compiler.source>
<maven.compiler.target>${target.java.version}</maven.compiler.target>
<log4j.version>2.14.1</log4j.version>
</properties>
<dependencies>
<!-- Core Flink dependencies — provided scope keeps them out of the job JAR. -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<!-- Connector dependencies — use compile scope to bundle them in the job JAR. -->
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-kafka</artifactId>
<version>${vvr.version}</version>
</dependency>
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-mysql</artifactId>
<version>${vvr.version}</version>
</dependency>
<!-- Logging framework — runtime scope excludes it from the application JAR. -->
<dependency>
<groupId>org.apache.logging.log4j</groupId>
<artifactId>log4j-slf4j-impl</artifactId>
<version>${log4j.version}</version>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.apache.logging.log4j</groupId>
<artifactId>log4j-api</artifactId>
<version>${log4j.version}</version>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.apache.logging.log4j</groupId>
<artifactId>log4j-core</artifactId>
<version>${log4j.version}</version>
<scope>runtime</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.11.0</version>
<configuration>
<source>${target.java.version}</source>
<target>${target.java.version}</target>
</configuration>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.2.0</version>
<executions>
<execution>
<phase>package</phase>
<goals>
<goal>shade</goal>
</goals>
<configuration>
<artifactSet>
<excludes>
<exclude>org.apache.flink:force-shading</exclude>
<exclude>com.google.code.findbugs:jsr305</exclude>
<exclude>org.slf4j:*</exclude>
<exclude>org.apache.logging.log4j:*</exclude>
</excludes>
</artifactSet>
<filters>
<filter>
<artifact>*:*</artifact>
<excludes>
<exclude>META-INF/*.SF</exclude>
<exclude>META-INF/*.DSA</exclude>
<exclude>META-INF/*.RSA</exclude>
</excludes>
</filter>
</filters>
<transformers>
<transformer
implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
<mainClass>com.aliyun.FlinkDemo</mainClass>
</transformer>
</transformers>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
Próximos passos
Conectores suportados — lista completa de conectores compatíveis com DataStream
Início rápido para jobs JAR do Flink — tutorial passo a passo ponta a ponta
Mapa de desenvolvimento de jobs e Desenvolver jobs Python — alternativas em SQL e Python