Este tópico descreve como gravar e consultar dados em tabelas Apache Hudi no E-MapReduce (EMR).
Pré-requisitos
Antes de começar, verifique se você tem:
Um cluster EMR executando o EMR V3.32.0 ou posterior
Configurar o ambiente
No EMR V3.32.0 e versões posteriores, as dependências do Hudi já estão integradas ao Spark, Hive e Presto. Não são necessárias dependências adicionais de runtime. Adicione a seguinte dependência Maven ao seu arquivo pom.xml:
<dependency>
<groupId>org.apache.hudi</groupId>
<artifactId>hudi-spark_2.11</artifactId>
<version>${hudi_version}</version>
<scope>provided</scope>
</dependency>
Para o Spark 3, usehudi-spark_2.12comoartifactIdem vez dehudi-spark_2.11.
A versão do Hudi incluída no cluster varia conforme a versão do EMR. Consulte a tabela abaixo para identificar o valor correto de ${hudi_version}:
|
Versão do Hudi |
Versão do EMR |
|
0.6.0 |
EMR V3.32.0–V3.35.0; EMR V4.5.0–V4.9.0; EMR V5.1.0 |
|
0.8.0 |
EMR V3.36.1–V3.37.1; EMR V5.2.1–V5.3.1 |
|
0.9.0 |
EMR V3.38.0–V3.38.3; EMR V5.4.0–V5.4.3 |
|
0.10.0 |
EMR V3.39.1–V3.40.0; EMR V4.10.0; EMR V5.5.0–V5.6.0 |
|
0.11.0 |
EMR V3.42.0; EMR V5.8.0 |
|
0.12.0 |
EMR V3.43.0–V3.44.1; EMR V5.9.0–V5.10.1 |
|
0.12.2 |
EMR V3.45.0–V3.46.1; EMR V5.11.0–V5.12.1 |
|
0.13.1 |
EMR V3.47.0–V3.48.0; EMR V5.13.0–V5.14.0 |
Gravar dados
Todos os exemplos de gravação compartilham um conjunto comum de opções do Hudi. Defina-as uma única vez e passe-as para cada operação de escrita:
val spark = SparkSession
.builder()
.master("local[*]")
.appName("hudi test")
.config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
.getOrCreate()
import spark.implicits._
// Sample dataset
val df = (for (i <- 0 until 10) yield (i, s"a$i", 30 + i * 0.2, 100 * i + 10000, s"p${i % 5}"))
.toDF("id", "name", "price", "version", "dt")
// Common options shared across all write operations
val hudiOptions = Map(
TABLE_NAME -> "hudi_test_0",
RECORDKEY_FIELD_OPT_KEY -> "id",
PRECOMBINE_FIELD_OPT_KEY -> "version",
KEYGENERATOR_CLASS_OPT_KEY -> classOf[SimpleKeyGenerator].getName,
HIVE_PARTITION_EXTRACTOR_CLASS_OPT_KEY -> classOf[MultiPartKeysValueExtractor].getCanonicalName,
PARTITIONPATH_FIELD_OPT_KEY -> "dt",
HIVE_PARTITION_FIELDS_OPT_KEY -> "ds",
META_SYNC_ENABLED_OPT_KEY -> "true",
HIVE_USE_JDBC_OPT_KEY -> "false",
HIVE_DATABASE_OPT_KEY -> "default",
HIVE_TABLE_OPT_KEY -> "hudi_test_0"
)
Principais parâmetros:
|
Parâmetro |
Descrição |
|
|
Nome da tabela Hudi. |
|
|
Campo usado como chave do registro. Registros com a mesma chave são deduplicados ou atualizados conforme o tipo de operação. |
|
|
Campo usado para pré-combinação antes da gravação. Quando dois registros têm a mesma chave, o Hudi mantém aquele com o maior valor neste campo. |
|
|
Classe geradora de chaves. A |
|
|
Campo do DataFrame usado como caminho de partição do Hudi. |
|
|
Nome da coluna de partição do Hive para sincronização de metadados. |
|
|
Ativa a sincronização de metadados do Hive. Obrigatório ao consultar a tabela com Hive ou Presto. |
|
|
Banco de dados Hive de destino para a sincronização de metadados. |
|
|
Nome da tabela Hive de destino para a sincronização de metadados. |
Inserir dados
df.write.format("hudi")
.options(hudiOptions)
.option(OPERATION_OPT_KEY, INSERT_OPERATION_OPT_VAL)
.option(INSERT_PARALLELISM, "8")
.option(UPSERT_PARALLELISM, "8")
.mode(Overwrite)
.save("/tmp/hudi/h0")
Para atualizar registros existentes, substitua INSERT_OPERATION_OPT_VAL por UPSERT_OPERATION_OPT_VAL.
Excluir dados
df.write.format("hudi")
.options(hudiOptions)
.option(OPERATION_OPT_KEY, DELETE_OPERATION_OPT_VAL)
.option(DELETE_PARALLELISM, "8")
.mode(Append)
.save("/tmp/hudi/h0")
A operação de exclusão localiza os registros pela chave definida em RECORDKEY_FIELD_OPT_KEY e os remove da tabela.
Consultar dados
O Hudi está integrado ao Spark, Hive e Presto no EMR. Portanto, não é preciso adicionar dependências extras para executar consultas.
Para consultar uma tabela Hudi com Hive ou Presto, ative a sincronização de metadados durante a gravação definindoMETA_SYNC_ENABLED_OPT_KEYcomo"true". Consulte Gravar dados .
Comportamento do formato de entrada
O Hudi do EMR e o Hudi open source tratam os formatos de entrada de maneiras distintas:
Hudi open source: Defina
hive.input.formatcomoorg.apache.hudi.hadoop.hive.HoodieCombineHiveInputFormattanto para tabelas Copy on Write quanto Merge on Read.Hudi do EMR: Não é preciso especificar um formato de entrada para tabelas Copy on Write, pois elas se adaptam automaticamente aos formatos de entrada do Hudi.