Este white paper descreve como usar a ferramenta de benchmark Nexmark para avaliar o desempenho de processamento de fluxo do Realtime Compute for Apache Flink.
Visão geral do desempenho
O Nexmark é um benchmark de desempenho padrão da indústria para motores de processamento de fluxo. Ele inclui 19 consultas padrão que abrangem cenários típicos, como filtragem, agregação, joins e janelas. Este documento usa a ferramenta de teste Nexmark para realizar uma avaliação completa de desempenho do Realtime Compute for Apache Flink com uma configuração de 8 CU e uma base de 100 milhões de registros de entrada para cada consulta. Os resultados dos testes mostram que:
Consultas simples, como q0, q1 e q2, atingem um RPS de 4 milhões a 6,5 milhões de registros por segundo.
Consultas complexas de agregação e janela, como q4, q5 e q16, alcançam um RPS entre 150.000 e 630.000 registros por segundo.
No geral, o Realtime Compute for Apache Flink entrega 3,24 vezes o desempenho Nexmark em comparação ao Flink open source.
Ferramenta de teste
O Nexmark é uma suíte de testes de benchmark de desempenho padrão para motores de processamento de fluxo. O modelo de teste é o seguinte:
Tabela de origem Nexmark: Gera dados de teste (eventos Person, Auction e Bid) em um TPS especificado.
Transformações: 19 consultas padrão Nexmark que cobrem cenários típicos, incluindo filtragem, transformação, agregação, joins e janelas.
Tabela de destino Blackhole: Grava dados em um destino blackhole para eliminar interferências de desempenho do armazenamento externo. Isso permite que a avaliação foque na capacidade de processamento do próprio motor Flink.
A ferramenta de teste Nexmark usada neste documento foi implementada com base na OpenAPI do Realtime Compute for Apache Flink. Ela automatiza todo o fluxo de trabalho, incluindo criação de jobs, implantação, monitoramento e coleta de resultados. Não é necessário escrever SQL manualmente ou criar jobs no console.
Ambiente de teste
Os jobs do Flink neste teste usaram as seguintes configurações de otimização:
|
Parâmetro |
Valor |
Descrição |
|
table.exec.mini-batch.enabled |
true |
Ativa a agregação Mini-Batch. |
|
table.exec.mini-batch.allow-latency |
2s |
Intervalo de buffer do Mini-Batch. |
|
table.optimizer.distinct-agg.split.enabled |
true |
Habilita a otimização de divisão para agregação Distinct. |
|
execution.checkpointing.interval |
3min |
Intervalo de checkpoint. |
Pré-requisitos
Java Development Kit (JDK) 1.8.x ou posterior instalado.
Realtime Compute for Apache Flink ativado e workspace criado. Para mais informações, consulte Ativar o Realtime Compute for Apache Flink.
AccessKey ID e AccessKey Secret da sua conta Alibaba Cloud obtidos.
Procedimento
Etapa 1: Baixe a ferramenta de teste
Baixe e extraia o pacote da ferramenta de teste Nexmark nexmark-flink.tar.gz.
A estrutura de diretórios após a extração é a seguinte:
nexmark-flink/
├── run_nexmark.sh # Test entry script
├── nexmark_env.sh # Environment variable configuration file (requires editing)
├── bin/ # Runtime scripts
├── conf/ # Flink job configurations
├── lib/ # JAR files (to be uploaded to the console)
└── queries-vvp/ # Nexmark Query SQL files
Etapa 2: Fazer upload do JAR Nexmark
Faça login no console do Realtime Compute for Apache Flink.
Clique em no espaço de projeto desejado. No painel de navegação à esquerda, selecione .
Selecione e faça upload do arquivo
nexmark-flink-0.2-SNAPSHOT.jar. O arquivo está localizado no diretórionexmark-flink/libda ferramenta de teste.-
Após concluir o upload, clique em no nome do arquivo para copiar seu endereço OSS. Você precisará desse endereço em uma etapa posterior de configuração. O formato do caminho do arquivo varia conforme o tipo de armazenamento:
-
Armazenamento em OSS Bucket:
oss://<nome do OSS Bucket>/artifacts/namespaces/<nome do espaço de projeto>/<nome do arquivo>Exemplo:
oss://oss-test/artifacts/namespaces/flink-default/nexmark-flink-0.2-SNAPSHOT.jar -
Armazenamento totalmente gerenciado:
oss://flink-fullymanaged-<ID do workspace>/artifacts/namespaces/<nome do espaço de projeto>/<nome do arquivo>Exemplo: oss://flink-fullymanaged-e6a123456789/artifacts/namespaces/flink-default/nexmark-flink-0.2-SNAPSHOT.jar
Para visualizar o tipo de armazenamento do seu workspace, acesse o console de gerenciamento do Realtime Compute for Apache Flink, localize o workspace desejado e clique em Details na coluna Actions.
-
Etapa 3: Configure parâmetros de execução
Edite o arquivo nexmark-flink/nexmark_env.sh e defina os seguintes parâmetros.
|
Parâmetro |
Descrição |
Exemplo |
|
END_POINT |
O endpoint do Realtime Compute for Apache Flink. Selecione o endpoint correspondente à sua região. Para mais informações, consulte Endpoints. |
ververica.cn-hangzhou.aliyuncs.com |
|
AK |
O AccessKey ID da sua conta Alibaba Cloud. |
- |
|
SK |
O AccessKey Secret da sua conta Alibaba Cloud. |
- |
|
WORK_SPACE |
O ID do seu workspace. |
e6a123456789 |
|
NAMESPACE |
O nome do seu espaço de projeto. |
flink-default |
|
NEXMARK_JAR |
O endereço OSS do arquivo JAR enviado na Etapa 2. |
oss://flink-fullymanaged-e6a123456789/artifacts/namespaces/flink-default/nexmark-flink-0.2-SNAPSHOT.jar |
|
FLINK_VERSION |
A versão do motor Flink a ser testada. |
vvr-11.5-jdk11-flink-1.20 |
|
QUERIES |
Especifique as consultas a serem executadas. Separe múltiplas consultas por vírgulas, por exemplo, |
all |
Executar todas as consultas consome tempo. Cada consulta passa por etapas como criação de job, geração de dados e execução de computação. Recomendamos executar primeiro uma única consulta (por exemplo, definir QUERIES como q0) para verificar se a configuração do ambiente e os parâmetros estão corretos antes de iniciar um teste em larga escala.
Etapa 4: Execute o teste
-
No diretório
nexmark-flink, execute o seguinte comando../run_nexmark.sh A ferramenta de teste cria e executa automaticamente os jobs Nexmark usando a OpenAPI.
-
Ao final do teste, a duração de cada consulta é exibida em milissegundos. O exemplo a seguir mostra uma saída de amostra:
INFO com.github.nexmark.flink.vvp.Nexmark - q0 13078 ============================================================================ ✓ Benchmark execution completed successfully ============================================================================
Resultados de desempenho
A tabela a seguir compara o desempenho Nexmark entre o Flink open source (1.20.4) e o Realtime Compute for Apache Flink (vvr-11.5-jdk11-flink-1.20) em uma configuração com 8 CUs. Cada consulta processa uma entrada de 100 milhões de registros. RPS = Número de registros de entrada ÷ Duração.
Os dados de teste a seguir foram coletados em um ambiente de hardware específico e com versões específicas do motor. O desempenho real pode variar devido a atualizações de hardware e do motor. Estes resultados servem apenas como referência.
|
Consulta |
Flink open source no ECS Versão: 1.20.4 |
Realtime Compute for Apache Flink Versão: vvr-11.5-jdk11-flink-1.20 |
|||
|
Duração (ms) |
RPS |
Duração (ms) |
RPS |
RPS vs. open source (×) |
|
|
q0 |
58848 |
1.699.293 |
23450 |
4.264.392 |
2,51 |
|
q1 |
57045 |
1.753.002 |
22824 |
4.381.353 |
2,50 |
|
q2 |
51890 |
1.927.154 |
15224 |
6.568.576 |
3,41 |
|
q3 |
84986 |
1.176.664 |
21558 |
4.638.649 |
3,94 |
|
q4 |
553426 |
180.693 |
157117 |
636.468 |
3,52 |
|
q5 |
365636 |
273.496 |
357547 |
279.684 |
1,02 |
|
q7 |
1257452 |
79.526 |
333837 |
299.547 |
3,77 |
|
q8 |
79788 |
1.253.321 |
29939 |
3.340.125 |
2,67 |
|
q9 |
2324518 |
43.020 |
266563 |
375.146 |
8,72 |
|
q10 |
189985 |
526.357 |
51202 |
1.953.049 |
3,71 |
|
q11 |
408384 |
244.868 |
145983 |
685.011 |
2,80 |
|
q12 |
121554 |
822.680 |
36991 |
2.703.360 |
3,29 |
|
q14 |
68903 |
1.451.316 |
20012 |
4.997.002 |
3,44 |
|
q15 |
183709 |
544.339 |
42734 |
2.340.057 |
4,30 |
|
q16 |
917597 |
108.980 |
337293 |
296.478 |
2,72 |
|
q17 |
102847 |
972.318 |
27076 |
3.693.308 |
3,80 |
|
q18 |
574949 |
173.928 |
96335 |
1.038.044 |
5,97 |
|
q19 |
586287 |
170.565 |
95121 |
1.051.293 |
6,16 |
|
q20 |
1340638 |
74.591 |
231482 |
431.999 |
5,79 |
|
q21 |
127089 |
786.850 |
39693 |
2.519.336 |
3,20 |
|
q22 |
94830 |
1.054.519 |
31228 |
3.202.254 |
3,04 |
|
Total |
2383209 |
49.695.131 |
9550361 |
15.317.480 |
3,24 |
Procedimento de teste do Flink open source
Para reproduzir os resultados no Flink open source em ECS, siga este procedimento.
Preparação do ambiente
Crie um cluster Flink usando EMR on ECS com a seguinte configuração:
Versão do EMR: EMR-5.21.0
Especificação de hardware: Três instâncias ecs.g6a.xlarge (4 vCPU / 16 GiB), incluindo um nó mestre e dois nós principais.
Ative os serviços Hadoop e HDFS.
-
Configure o login sem senha entre todos os nós. Por exemplo, faça upload do seu arquivo de chave privada (como
key.pem) para o nó mestre. Em seguida, adicione a seguinte configuração ao arquivo~/.ssh/configno nó mestre. Substitua os endereços IP e o caminho do arquivo pelos seus valores reais.Host 192.168.0.0 HostName 192.168.0.0 User root IdentityFile /path/to/key.pem StrictHostKeyChecking no Host 192.168.0.1 HostName 192.168.0.1 User root IdentityFile /path/to/key.pem StrictHostKeyChecking no Host 192.168.0.2 HostName 192.168.0.2 User root IdentityFile /path/to/key.pem StrictHostKeyChecking noUse ssh para verificar se o login sem senha entre os nós está funcionando corretamente. Se ocorrer um erro de "bad permissions", execute
chmod 600 /path/to/key.pempara corrigir as permissões.
Preparação do software
-
Baixe o pacote Flink desejado (Apache Flink Downloads) e o pacote de teste Nexmark (nexmark-flink.tgz). Faça upload dos pacotes para o nó mestre e extraia-os.
tar -zxvf flink-1.20.4-bin-scala_2.12.tgz tar -zxvf nexmark-flink.tgz mv flink-1.20.4 flink mv nexmark-flink nexmark
-
Copie os arquivos JAR do diretório
nexmark/libparaflink/lib. Esses JARs contêm o gerador de dados Nexmark.cp nexmark/lib/* flink/lib/
-
Defina as variáveis de ambiente. Edite
~/.bashrc, adicione a seguinte configuração e executesource ~/.bashrcpara que as alterações entrem em vigor.Configure os caminhos com base no seu ambiente real.
export JAVA_HOME=/etc/alternatives/java_sdk_11 export PATH=$JAVA_HOME/bin:$PATHexport FLINK_HOME=/mnt/disk1/flink export HADOOP_CLASSPATH=$(/opt/apps/HADOOP-COMMON/hadoop-common-current/bin/hadoop classpath)
Configuração e inicialização do cluster
-
Configure os workers do Flink. Este teste usa oito TaskManagers, implantados da seguinte forma: dois no nó mestre e três em cada nó principal.
Edite
flink/conf/workers, certificando-se de substituir o endereço IP pelo valor real.192.168.0.0 192.168.0.0 192.168.0.1 192.168.0.1 192.168.0.1 192.168.0.2 192.168.0.2 192.168.0.2
-
Substitua
flink/conf/config.yamlpornexmark/conf/config.yamle atualize os seguintes itens de configuração:jobmanager.rpc.address: O endereço IP do nó mestre, como192.168.0.0.state.checkpoints.dir: O caminho HDFS, comohdfs:///checkpointstaskmanager.memory.process.size:4G
Edite
nexmark/conf/nexmark.yamle definanexmark.metric.reporter.hostcomo o endereço IP do nó mestre.-
Distribua os diretórios
flinkenexmark, bem como as configurações de variáveis de ambiente, para cada nó principal.Substitua os endereços IP pelos seus valores reais.
scp -r flink 192.168.0.1:/mnt/disk1/ scp -r flink 192.168.0.2:/mnt/disk1/ scp -r nexmark 192.168.0.1:/mnt/disk1/ scp -r nexmark 192.168.0.2:/mnt/disk1/ scp ~/.bashrc 192.168.0.1:~/ scp ~/.bashrc 192.168.0.2:~/Após concluir a distribuição, execute
source ~/.bashrcem cada nó principal para que as variáveis de ambiente entrem em vigor.
-
No nó mestre, inicie o cluster Flink.
flink/bin/start-cluster.sh
-
Inicialize o ambiente de teste Nexmark. Este script configura o Metric Reporter necessário em cada nó.
nexmark/bin/setup_cluster.sh
Defina limites de recursos
Após a inicialização do cluster Flink, use cgroups para limitar o uso de CPU de cada processo TaskManager a 75%. Isso evita que os TaskManagers percam conexão devido a timeouts de heartbeat causados por contenção de recursos.
Execute os seguintes comandos em todos os nós onde os TaskManagers estão em execução, incluindo o nó mestre:
yum install -y libcgroup libcgroup-tools
cgcreate -t root:root -a root:root -g cpu,memory:mygroup
echo 100000 > /sys/fs/cgroup/cpu/mygroup/cpu.cfs_period_us
echo 300000 > /sys/fs/cgroup/cpu/mygroup/cpu.cfs_quota_us
echo $((12 * 1024 * 1024 * 1024)) > /sys/fs/cgroup/memory/mygroup/memory.limit_in_bytes
jps | grep TaskManagerRunner | awk '{print $1}' | xargs cgclassify -g cpu,memory:mygroup
Aqui, cpu.cfs_quota_us / cpu.cfs_period_us = 300000 / 100000 = 3, o que significa que o cgroup pode usar no máximo 3 núcleos de CPU (75% de 4 vCPUs).
No ambiente totalmente gerenciado do Realtime Compute for Apache Flink, não é necessário definir manualmente um limite de uso de recursos. O serviço utiliza integralmente os recursos de computação adquiridos.
Execute o Nexmark
No nó mestre, execute o seguinte comando e aguarde a conclusão para visualizar os resultados:
nexmark/bin/run_query.sh q0,q1,q2,q3,q4,q5,q7,q8,q9,q10,q11,q12,q14,q15,q16,q17,q18,q19,q20,q21,q22