Uso de un trabajo de Flink de DLI para sincronizar datos de Kafka con un clúster de DWS en tiempo real
Esta práctica demuestra cómo utilizar DLI Flink (se utiliza Flink 1.15 como ejemplo) para sincronizar los datos de consumo de Kafka con DWS en tiempo real. El proceso de demostración incluye la escritura y actualización de datos existentes en tiempo real.
- Para obtener más información, consulte ¿Qué es Data Lake Insight?
- Para obtener más información sobre Kafka, consulte ¿Qué es DMS for Kafka?
Esta práctica dura unos 90 minutos. Los servicios en la nube utilizados en esta práctica incluyen Virtual Private Cloud (VPC) y subredes, Elastic Load Balance (ELB), Elastic Cloud Server (ECS), Object Storage Service (OBS), Distributed Message Service (DMS) for Kafka, Data Lake Insight (DLI) y Data Warehouse Service (DWS). El proceso básico es el siguiente:
- Preparativos: Registre una cuenta y prepare la red.
- Paso 1: Crear una instancia de Kafka: Compre DMS for Kafka y prepare los datos de Kafka.
- Paso 2: Crear un clúster de DWS y una tabla de destino: Compre un clúster de DWS y vincúlelo a una EIP.
- Paso 3: Crear un grupo de recursos elásticos de DLI y una cola: Cree un grupo de recursos elásticos de DLI y agregue colas al grupo de recursos.
- Paso 4: Crear una conexión de origen de datos mejorada para Kafka y DWS: Conecte Kafka y DWS.
- Paso 5: Preparar la herramienta dws-connector-flink para interconectar DWS con Flink: Utilice este complemento para importar datos de MySQL a DWS de manera eficiente.
- Paso 6: Crear y editar un trabajo de Flink de DLI: Cree un trabajo de Flink SQL y configure el código SQL.
- Paso 7: Crear y modificar mensajes en el cliente de Kafka: Importe datos a DWS en tiempo real.
Descripción del escenario
Supongamos que los datos de muestra de la fuente de datos Kafka son una tabla de información de usuario, como se muestra en Tabla 1, que contiene los campos id, name y age. El campo id es único y fijo, que es compartido por múltiples sistemas de servicio. Por lo general, no es necesario modificar el campo id. Solo es necesario modificar los campos name y age.
Utilice Kafka para generar los siguientes tres grupos de datos y utilice trabajos de Flink de DLI para sincronizar los datos con DWS: Cambie los usuarios cuyos ID son 2 y 3 a jim y tom y utilice trabajos de Flink de DLI para actualizar los datos y sincronizarlos con DWS.
Restricciones
- Asegúrese de que VPC, el ECS, el OBS, Kafka, el DLI y el DWS estén en la misma región, por ejemplo, CN-Hong Kong.
- Asegúrese de que Kafka, el DLI y el DWS puedan comunicarse entre sí. En esta práctica, Kafka y el DWS se crean en la misma región y VPC, y los grupos de seguridad de Kafka y el DWS permiten el segmento de red de las colas de DLI.
- Para garantizar que el enlace entre DLI y DWS sea estable, vincule el servicio ELB al clúster de DWS creado.
Preparativos
- Ha creado un ID de Huawei y habilitado los servicios de Huawei Cloud. La cuenta no puede estar en mora o congelada.
- Ha creado una VPC y una subred. Para obtener más información, consulte Creación de una VPC.
Paso 1: Crear una instancia de Kafka
- Inicie sesión en la consola de Kafka.
- Configure los siguientes parámetros según se indique. Mantenga la configuración predeterminada para otros parámetros.
Tabla 2 Parámetros de instancia de Kafka Parámetro
Valor
Billing Mode
Pago por uso
Region
CN-Hong Kong
AZ
AZ 1 (Si no está disponible, seleccione otra AZ).
Bundle
Inicial
VPC
Seleccione una VPC creada. Si no hay ninguna VPC disponible, cree una.
Security Group
Seleccione un grupo de seguridad creado. Si no hay ningún grupo de seguridad disponible, cree uno.
Other parameters
Conserve el valor predeterminado.
- Haga clic en Confirm, confirme la información y haga clic en Submit. Espere hasta que la creación sea exitosa.
- En la lista de instancias de Kafka, haga clic en el nombre de la instancia de Kafka creada. Se muestra la página Basic Information.
- Elija Topics a la izquierda y haga clic en Create Topic.
Establezca Topic Name en topic-demo y conserve los valores predeterminados para los otros parámetros.
Figura 2 Creación de un tema
- Haga clic en OK. En la lista de temas, puede ver que topic-demo se ha creado correctamente.
- Elija Consumer Groups a la izquierda y haga clic en Create Consumer Group.
- Ingrese kafka01 para Consumer Group Name y haga clic en OK.
Paso 2: Crear un clúster de DWS y una tabla de destino
- Cree un balanceador de carga dedicado, establezca Network Type en IPv4 private network. Establezca Region y VPC en los mismos valores que los de la instancia de Kafka. En este ejemplo, establezca Region en CN-Hong Kong.
- Cree un clúster y vincule un balanceador de carga elástico al clúster. Para garantizar la conectividad de red, la región y la VPC del clúster de DWS deben ser las mismas que las de la instancia de Kafka. En esta práctica, la región y la VPC son CN-Hong Kong. La VPC debe ser la misma que la creada para Kafka.
- Inicie sesión en la consola de DWS. Seleccione Dedicated Clusters > Clusters. Busque el clúster de destino y haga clic en Log In en la columna Operation.
Este modo de inicio de sesión solo está disponible para clústeres de la versión 8.1.3.x. Para los clústeres de la versión 8.1.2 o anterior, debe usar gsql para iniciar sesión.
- Una vez que el inicio de sesión se realiza correctamente, se muestra el editor SQL.
- Copie la siguiente sentencia SQL. En la ventana SQL, haga clic en Execute SQL para crear la tabla de destino user_dws.
1 2 3 4 5 6
CREATE TABLE user_dws ( id int, name varchar(50), age int, PRIMARY KEY (id) );
Paso 3: Crear un grupo de recursos elásticos de DLI y una cola
- Inicie sesión en la consola de DLI.
- En el panel de navegación de la izquierda, seleccione Resources > Resource Pool.
- Haga clic en Buy Resource Pool en la esquina superior derecha, configure los siguientes parámetros según se indique y conserve la configuración predeterminada para otros parámetros.
Tabla 3 Parámetros Parámetro
Valor
Billing Mode
Pago por uso
Region
CN-Hong Kong
Name
dli_dws
Specifications
Estándar
CIDR Block
172.16.0.0/18. Debe estar en un segmento de red diferente al de Kafka y DWS. Por ejemplo, si Kafka y DWS están en el segmento de red 192.168.x.x, seleccione 172.16.x.x para DLI.
- Haga clic en Buy y haga clic en Submit.
Una vez creado el grupo de recursos, siga con el siguiente paso.
- En la página de grupo de recursos elásticos, busque la fila que contiene el grupo de recursos creado, haga clic en Add Queue en la columna Operation y configure los siguientes parámetros. Conserve los valores predeterminados para otros parámetros que no se describen en la tabla.
Tabla 4 Adición de una cola Parámetro
Valor
Name
dli_dws
Type
Cola de uso general
- Haga clic en Next y en OK. La cola se ha creado.
Paso 4: Crear una conexión de origen de datos mejorada para Kafka y DWS
- Actualice los permisos de la delegación de DLI.
- Vuelva a la consola de DLI y seleccione Global Configuration > Service Authorization a la izquierda.
- Seleccione DLI UserInfo Agency Access, DLI Datasource Connections Agency Access y DLI Notification Agency Access.
- Haga clic en Update y luego en OK.
Figura 3 Actualización de permisos de delegación de DLI
- En el grupo de seguridad de Kafka, permita el segmento de red donde se encuentra la cola de DLI.
- Vuelva a la consola de Kafka y haga clic en el nombre de la instancia de Kafka para ir a la página Basic Information. Vea el valor de Instance Address (Private Network) en la información de conexión y registre la dirección para su uso futuro. Figura 4 Dirección de red privada de Kafka
- Haga clic en el nombre del grupo de seguridad. Figura 5 Grupo de seguridad de Kafka
- Elija Inbound Rules > Add Rule, como se muestra en la siguiente figura. Agregue el segmento de red de la cola de DLI. En este ejemplo, el segmento de red es 172.16.0.0/18. Asegúrese de que el segmento de red sea el mismo que se ingresó durante Paso 3: Crear un grupo de recursos elásticos de DLI y una cola. Figura 6 Adición de reglas al grupo de seguridad de Kafka
- Haga clic en OK.
- Vuelva a la consola de Kafka y haga clic en el nombre de la instancia de Kafka para ir a la página Basic Information. Vea el valor de Instance Address (Private Network) en la información de conexión y registre la dirección para su uso futuro.
- Vuelva a la consola de gestión de DLI, haga clic en Datasource Connections a la izquierda, seleccione Enhanced y haga clic en Create.
- Configure los siguientes parámetros. Conserve los valores predeterminados para otros parámetros que no se describen en la tabla.
Tabla 5 Conexión de DLI a Kafka Parámetro
Valor
Connection Name
dli_kafka
Resource Pool
Seleccione la cola de DLI creada dli_dws.
VPC
Seleccione la VPC de Kafka.
Subnet
Seleccione la subred donde se encuentra Kafka.
Other parameters
Conserve el valor predeterminado.
Figura 7 Creación de una conexión mejorada
- Haga clic en OK. Espere hasta que la conexión de Kafka se haya creado correctamente.
Si no selecciona un grupo de recursos al crear una conexión de origen de datos mejorada, puede vincular uno manualmente después.
- Localice la fila que contiene la conexión de origen de datos de destino y elija More > Bind Resource Pool en la columna Operation.
- Seleccione un grupo de recursos y haga clic en OK.

- Elija Resources > Queue Management a la izquierda y elija More > Test Address Connectivity a la derecha de dli_dws.
- En el cuadro de dirección, ingrese la dirección IP privada y el número de puerto de la instancia de Kafka obtenida en 2.a. (Hay tres direcciones de Kafka. Ingrese solo una de ellas). Figura 8 Prueba de conectividad de Kafka
- Haga clic en Test para verificar que DLI se haya conectado correctamente a Kafka.
- Si la conexión falla, realice las siguientes operaciones:
- Inicie sesión en la consola de DWS, elija Dedicated Clusters > Clusters a la izquierda y haga clic en el nombre del clúster para ir a la página de detalles.
- Registre el nombre de dominio de la red privada, el número de puerto y la dirección de ELB del clúster de DWS para su uso futuro. Figura 9 Nombre de dominio privado y dirección de ELB
- Haga clic en el nombre del grupo de seguridad. Figura 10 Grupo de seguridad de DWS
- Elija Inbound Rules > Add Rule, como se muestra en la siguiente figura. Agregue el segmento de red de la cola de DLI. En este ejemplo, el segmento de red es 172.16.0.0/18. Asegúrese de que el segmento de red sea el mismo que se ingresó durante Paso 3: Crear un grupo de recursos elásticos de DLI y una cola. Figura 11 Adición de una regla al grupo de seguridad de DWS
- Haga clic en OK.
- Cambie a la consola de DLI, elija Resources > Queue Management a la izquierda y haga clic en More > Test Address Connectivity a la derecha de dli_dws.
- En el cuadro de direcciones, ingrese la dirección de ELB y el número de puerto del clúster de DWS obtenido en 11. Figura 12 Prueba de conectividad de DWS
- Haga clic en Test para verificar que DLI se haya conectado correctamente a DWS.
Paso 5: Preparar la herramienta dws-connector-flink para interconectar DWS con Flink
dws-connector-flink es una herramienta para interconectarse con Flink basada en las API JDBC de DWS. Durante la configuración del trabajo de DLI, esta herramienta y sus dependencias se almacenan en el directorio de carga de clases de Flink para mejorar la capacidad de importar trabajos de Flink a DWS.
- Vaya a https://mvnrepository.com/artifact/com.huaweicloud.dws usando un navegador.
- En la lista de software, seleccione Flink 1.15. En esta práctica, se selecciona DWS Connector Flink SQL 1 15.

- Seleccione la rama más reciente. La rama real está sujeta a la nueva rama publicada en el sitio web oficial.

- Haga clic en el ícono jar para descargar el archivo.

- Cree un bucket de OBS. En esta práctica, configure el nombre de bucket como obs-flink-dws y cargue el archivo dws-connector-flink-sql-1.15-2.12_2.0.0.r4.jar en el bucket de OBS. Asegúrese de que el bucket se encuentre en la misma región que DLI. En esta práctica, se utiliza la región CN-Hong Kong.
Paso 6: Crear y editar un trabajo de Flink de DLI
- Cree una política de delegación de OBS.
- Pase el cursor sobre el nombre de la cuenta en el extremo superior derecho de la consola y haga clic en Identity and Access Management.
- En el panel de navegación de la izquierda, seleccione Agencies y, a continuación, haga clic en Create Agency en el extremo superior derecho.
- Agency Name: dli_ac_obs
- Agency Type: Cloud service
- Cloud Service: Data Lake Insight (DLI)
- Validity Period: Unlimited
Figura 13 Creación de una delegación de OBS
- Haga clic en OK y, a continuación, haga clic en Authorize.
- En la página que se muestra, haga clic en Create Policy.
- Configure la información de la política. Ingrese un nombre de política, por ejemplo, dli_ac_obs, y seleccione JSON.
- En el área Policy Content, pegue una política personalizada. Reemplace OBS bucket name con el nombre real de bucket.
{ "Version": "1.1", "Statement": [ { "Effect": "Allow", "Action": [ "obs:object:GetObject", "obs:object:DeleteObjectVersion", "obs:bucket:GetBucketLocation", "obs:bucket:GetLifecycleConfiguration", "obs:object:AbortMultipartUpload", "obs:object:DeleteObject", "obs:bucket:GetBucketLogging", "obs:bucket:HeadBucket", "obs:object:PutObject", "obs:object:GetObjectVersionAcl", "obs:bucket:GetBucketAcl", "obs:bucket:GetBucketVersioning", "obs:bucket:GetBucketStoragePolicy", "obs:bucket:ListBucketMultipartUploads", "obs:object:ListMultipartUploadParts", "obs:bucket:ListBucketVersions", "obs:bucket:ListBucket", "obs:object:GetObjectVersion", "obs:object:GetObjectAcl", "obs:bucket:GetBucketPolicy", "obs:bucket:GetBucketStorage" ], "Resource": [ "OBS:*:*:object:*", "OBS:*:*:bucket:OBS bucket name" ] }, { "Effect": "Allow", "Action": [ "obs:bucket:ListAllMyBuckets" ] } ] } - Haga clic en Next.
- Seleccione la política personalizada creada.
- Haga clic en Next. Seleccione All resources.
- Haga clic en OK.
La autorización tarda entre 15 y 30 minutos en surtir efecto.
- Vuelva a la consola de gestión de DLI, seleccione Job Management > Flink Jobs a la izquierda y haga clic en Create Job en la esquina superior derecha.
- Establezca Type en Flink OpenSource SQL y Name en kafka-dws. Figura 14 Creación de un trabajo
- Haga clic en OK. Se mostrará la página para editar el trabajo.
- Establezca los siguientes parámetros a la derecha de la página. Conserve los valores predeterminados para otros parámetros que no se describen en la tabla.
Tabla 6 Parámetros de trabajo de Flink Parámetro
Valor
Queue
dli_dws
Flink Version
1.15
UDF Jar
Seleccione el archivo JAR en el bucket de OBS creado en Paso 5: Preparar la herramienta dws-connector-flink para interconectar DWS con Flink.
Agency
Seleccione la delegación creada en 1.
OBS Bucket
Seleccione el bucket creado en Paso 5: Preparar la herramienta dws-connector-flink para interconectar DWS con Flink.
Enable Checkpointing
Marque la casilla.
Other parameters
Conserve el valor predeterminado.
- Copie el siguiente código SQL en la ventana de código SQL de la izquierda.
Obtenga la dirección IP privada y el número de puerto de la instancia de Kafka de 2.a y obtenga el nombre de dominio privado de 11.
A continuación se describen los parámetros comunes del código de trabajo de Flink SQL:
- connector: tipo de conector del origen de datos. Para Kafka, establezca este parámetro en kafka. Para DWS, establezca este parámetro en dws. Para obtener más información, consulte Conectores.
- write.mode: modo de importación. El valor puede ser auto, copy_merge, copy_upsert, upsert, update, copy_update o copy_auto.
- autoFlushBatchSize: número de registros que se almacenarán en búfer antes de una descarga automática. Cuando el número de registros en búfer alcanza este valor, Flink escribe datos en el sistema de destino por lotes. Por ejemplo, Flink escribe datos en el sistema de destino Sink después de que se almacenen en búfer 5000 líneas de registros. En Flink SQL, el receptor (Sink) es el destino final de salida en el pipeline de procesamiento del flujo de datos. Es responsable de escribir los datos procesados en sistemas de almacenamiento externo o de enviarlos a aplicaciones descendentes.
- autoFlushMaxInterval: intervalo máximo entre descargas automáticas. Incluso si el número de registros en el búfer no alcanza el valor de autoFlushBatchSize, se activa una descarga después de que caduque el intervalo.
- key-by-before-sink: Si se deben agrupar los datos por clave principal antes de que se escriban en el receptor. Esto garantiza que los registros con la misma clave principal se escriban consecutivamente y es útil para algunos sistemas (como algunas bases de datos) que requieren que las claves principales estén ordenadas. Este parámetro tiene como objetivo resolver el problema de interbloqueo entre dos subtareas cuando adquieren bloqueos de filas basados en la clave principal de DWS, se producen múltiples escrituras simultáneas y write.mode es upsert. Esto sucede cuando un lote de datos escritos en el receptor por varias subtareas tiene más de un registro con la misma clave principal y el orden de estos registros con la misma clave principal es inconsistente.
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31
CREATE TABLE user_kafka ( id string, name string, age int ) WITH ( 'connector' = 'kafka', 'topic' = 'topic-demo', 'properties.bootstrap.servers' ='Private IP address and port number of the Kafka instance', 'properties.group.id' = 'kafka01', 'scan.startup.mode' = 'latest-offset', 'format' = 'json' ); CREATE TABLE user_dws ( id string, name string, age int, PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'dws', 'url'='jdbc:postgresql://DWS private network domain name:8000/gaussdb', 'tableName' = 'public.user_dws', 'username' = 'dbadmin', 'password' ='Password of database user dbadmin' 'writeMode' = 'auto', 'autoFlushBatchSize'='50000', 'autoFlushMaxInterval'='5s', 'key-by-before-sink'='true' ); INSERT INTO user_dws select * from user_kafka; -- Write the processing result to the sink.
- Haga clic en Format y en Save.
Haga clic en Format para formatear el código SQL. De lo contrario, se podría introducir caracteres nulos nuevos al copiar y pegar el código, lo que provocaría fallas en la ejecución del trabajo.
Figura 15 Sentencia SQL de un trabajo
- Vuelva a la página de inicio de la consola de DLI y seleccione Job Management > Flink Jobs a la izquierda.
- Haga clic en Start a la derecha del nombre del trabajo kafka-dws y en Start Now.
Espere aproximadamente 1 minuto y actualice la página. Si el estado es Running, el trabajo se ejecutó correctamente.
Figura 16 Estado de ejecución del trabajo
Paso 7: Crear y modificar mensajes en el cliente de Kafka
- Cree un ECS consultando el documento de ECS. Asegúrese de que la región y la VPC del ECS sean las mismas que las de Kafka.
- Instale JDK.
- Inicie sesión en el ECS, vaya al directorio /usr/local y descargue el paquete JDK.
1 2
cd /usr/local wget https://download.oracle.com/java/17/latest/jdk-17_linux-x64_bin.tar.gz
- Descomprima el paquete JDK descargado.
1tar -zxvf jdk-17_linux-x64_bin.tar.gz
- Ejecute el siguiente comando para abrir el archivo /etc/profile:
1vim /etc/profile - Presione i para ingresar al modo de edición y agregue el siguiente contenido al final del archivo /etc/profile:
1 2 3 4 5
export JAVA_HOME=/usr/local/jdk-17.0.7 #JDK installation directory export JRE_HOME=${JAVA_HOME}/jre export CLASSPATH=.:${JAVA_HOME}/lib:${JRE_HOME}/lib:${JAVA_HOME}/test:${JAVA_HOME}/lib/gsjdbc4.jar:${JAVA_HOME}/lib/dt.jar:${JAVA_HOME}/lib/tools.jar:$CLASSPATH export JAVA_PATH=${JAVA_HOME}/bin:${JRE_HOME}/bin export PATH=$PATH:${JAVA_PATH}

- Presione Esc e ingrese :wq! para guardar la configuración y salir.
- Ejecute el siguiente comando para que las variables de entorno surtan efecto:
1source /etc/profile
- Ejecute el siguiente comando. Si se muestra la siguiente información, el JDK se ha instalado correctamente:
1java -version
- Inicie sesión en el ECS, vaya al directorio /usr/local y descargue el paquete JDK.
- Instale el cliente Kafka.
- Vaya al directorio /opt y ejecute el siguiente comando para obtener el paquete de software del cliente Kafka.
1 2
cd /opt wget https://archive.apache.org/dist/kafka/2.7.2/kafka_2.12-2.7.2.tgz
- Descomprima el paquete de software descargado.
1tar -zxf kafka_2.12-2.7.2.tgz
- Vaya al directorio del cliente Kafka.
1cd /opt/kafka_2.12-2.7.2/bin
- Vaya al directorio /opt y ejecute el siguiente comando para obtener el paquete de software del cliente Kafka.
- Ejecute el siguiente comando para conectarse a Kafka: {Connection address} indica la dirección de conexión de red interna de Kafka. Para obtener más detalles sobre cómo obtener la dirección, consulte 2.a. topic indica el nombre del tema de Kafka creado en 5.
1./kafka-console-producer.sh --broker-list {connection address} --topic {Topic name}
A continuación se muestra un ejemplo:
./kafka-console-producer.sh --broker-list 192.168.0.136:9092,192.168.0.214:9092,192.168.0.217:9092 --topic topic-demo

Si se muestra > y no se muestra ningún otro mensaje de error, la conexión se realiza correctamente.
- En la ventana del cliente Kafka conectado, copie el siguiente contenido (una línea a la vez) según los datos planificados en el Descripción del escenario y presione Enter para generar mensajes:
1 2 3
{"id":"1","name":"lily","age":"16"} {"id":"2","name":"lucy","age":"17"} {"id":"3","name":"lilei","age":"15"}

- Vuelva a la consola de DWS, seleccione Dedicated Clusters > Clusters a la izquierda y haga clic en Log In a la derecha del clúster de DWS. Se muestra la página SQL.
- Ejecute la siguiente sentencia SQL para verificar que los datos se importen correctamente a la base de datos en tiempo real:
1SELECT * FROM user_dws ORDER BY id;

- Vuelva a la ventana del cliente para conectarse a Kafka en el ECS, copie el siguiente contenido (una línea a la vez) y presione Enter para generar mensajes.
1 2
{"id":"2","name":"jim","age":"17"} {"id":"3","name":"tom","age":"15"}
- Vuelva a la ventana SQL abierta de DWS y ejecute la siguiente sentencia SQL. Se encuentra que los nombres cuyos ID son 2 y 3 han cambiado a jim y tom. La descripción del escenario es la esperada. Fin de esta práctica.
1SELECT * FROM user_dws ORDER BY id;
