DataWorks propose des vérifications intégrées, telles que la revue de code avant le déploiement des nœuds et les contrôles dans le Data Governance Center. DataWorks prend également en charge les vérifications personnalisées. Développez un programme de vérification personnalisé selon vos besoins métier et connectez-le à DataWorks pour gérer les processus des nœuds. Cette rubrique explique comment s'abonner aux événements de changement d'état d'une instance de nœud planifié à l'aide du module OpenEvent de la plateforme ouverte DataWorks. Les changements d'état d'une instance de nœud planifié sont obtenus depuis Operation Center.
Contexte
Pour plus d'informations sur les fonctionnalités et les concepts de base de la plateforme ouverte DataWorks abordés dans cette rubrique, consultez Présentation.
Description de la configuration de l'abonnement
Lorsque vous configurez le paramètre Pattern Content pour une règle d'événement dans EventBridge, définissez le type sur dataworks:InstanceStatusChanges:InstanceStatusChanges. Vous pouvez ainsi vous abonner aux événements de changement d'état d'une instance de nœud planifié.
Prérequis
EventBridge est activé. Pour plus d'informations, consultez Facturation.
DataWorks est activé. Pour plus d'informations, consultez Guide d'achat.
Un espace de travail a été créé dans DataWorks. Pour plus d'informations, consultez Créer un espace de travail.
Procédure
Étape 1 : Configurer un bus personnalisé
Cette section décrit les étapes de configuration principales et les précautions à prendre pour configurer un bus personnalisé. Pour plus d'informations sur l'activation et la configuration de l'abonnement aux messages d'événement, consultez Activer l'abonnement aux messages.
Connectez-vous à la console EventBridge. Dans le volet de navigation de gauche, cliquez sur Event Buses.
Dans le coin supérieur droit de la page Event Buses, cliquez sur Quickly Create pour créer un bus personnalisé.
-
Dans le volet de navigation de gauche, cliquez sur Event Buses. Sur la page Event Buses, localisez le bus personnalisé que vous avez créé et cliquez sur son nom. La page Overview du bus personnalisé s'affiche.
Dans le volet de navigation de gauche, cliquez sur Event Rules. Sur la page qui s'affiche, cliquez sur Create Rule pour créer une règle d'événement.
-
Dans cet exemple, le bus personnalisé est configuré pour recevoir les messages d'événement de validation de nœud et les messages d'événement de déploiement de nœud. Le contenu suivant fournit un exemple de configuration d'une démo et des paramètres principaux pour la règle d'événement :
Configure Basic Info: À cette étape, configurez le paramètre Name.
-
Configure Event Pattern :
Event Source Type : Définissez-le sur Custom Event Source.
Event Source : Laissez ce paramètre vide.
-
Pattern Content : Configurez ce paramètre au format JSON. Saisissez le contenu suivant.
{ "source": [ "acs.dataworks" ], "type": [ "dataworks:InstanceStatusChanges:InstanceStatusChanges" ] }source : l'identifiant du service dans lequel un événement se produit. Définissez ce paramètre sur acs.dataworks.
type : le type de l'événement qui se produit dans le service. Définissez ce paramètre sur dataworks:InstanceStatusChanges:InstanceStatusChanges.
Event Pattern Debugging : Modifiez les valeurs des paramètres source et type dans cette section, puis cliquez sur Test. Si le test réussit, cliquez sur Next Step.

-
Configure Targets :
Service Type : Définissez-le sur HTTPS ou HTTP. Pour plus d'informations sur les types de service, consultez Gérer les règles d'événement.
URL : Saisissez l'URL pour recevoir les messages poussés par le bus personnalisé, par exemple
https://Server address:Port number/event/consumer.Body : Définissez-le sur Complete Event.
Network Type : Définissez-le sur Internet.

Étape 2 : Configurer un canal de distribution d'événements
-
Connectez-vous à la console DataWorks. Dans la région cible, cliquez sur dans le volet de navigation de gauche. Cliquez sur Go to Open Platform pour ouvrir la page Developer Backend.
-
Sur la page Developer Backend, cliquez sur OpenEvent dans le volet de navigation de gauche. Sur la page qui s'affiche, cliquez sur Add Event Distribution Channel et configurez les paramètres dans la boîte de dialogue.
Workspace for Event Distribution : Sélectionnez un espace de travail que vous avez créé.
Custom Bus in Eventbridge for Distribution : Sélectionnez le bus d'événements créé à l'étape 1.
Après avoir enregistré le canal de distribution d'événements, dans la colonne Operation du canal, cliquez sur Enable pour activer le nouveau canal de distribution d'événements.
Développer un programme de service
Exemple de code
Modifiez le code pour obtenir les messages qu'EventBridge pousse vers le service et générer les messages.
package com.aliyun.dataworks.demo;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.aliyun.dataworks.config.Constants;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
/**
* @author dataworks demo
*/
@RestController
@RequestMapping("/event")
public class ExtensionsController {
/**
* Receive event messages that are sent from EventBridge.
* @param jsonParam
*/
@PostMapping("/consumer")
public void consumerEventBridge(@RequestBody String jsonParam){
JSONObject jsonObj = JSON.parseObject(jsonParam);
String eventCode = jsonObj.getString(Constants.EVENT_CODE_FILED);
if(Constants.INSTANCE_STATUS_EVENT_CODE.equals(eventCode)){
JSONObject dataParam = JSON.parseObject(jsonObj.getString("data"));
// The time when the auto triggered node instance started to wait for the scheduling time.
System.out.println("beginWaitTimeTime: "+ dataParam.getString("beginWaitTimeTime"));
//DagId
System.out.println("dagId: "+ dataParam.getString("dagId"));
// The type of the directed acyclic graph (DAG). Valid values:
// 0: for auto triggered nodes
// 1: for manually triggered nodes
// 2: for smoke testing
// 3: for nodes for which you backfill data
// 4: for manually triggered workflows
// 5: for temporary workflows
System.out.println("dagType: "+dataParam.getString("dagType"));
// The type of the node. Valid values:
// NORMAL(0): The node is an auto triggered node. The scheduling system regularly runs the node.
// MANUAL(1): The node is a manually triggered node. The scheduling system does not regularly run the node.
// PAUSE(2): The node is a frozen node. The scheduling system regularly runs the node but sets the node status to Failed when the scheduling system starts to run the node.
// SKIP(3): The node is a dry-run node. The scheduling system regularly runs the node but sets the node status to Succeeded when the scheduling system starts to run the node.
// SKIP_UNCHOOSE(4): The node is an unselected node in a temporary workflow. This type of node exists only in temporary workflows. The scheduling system sets the node status to Succeeded when the scheduling system starts to run the node.
// SKIP_CYCLE(5): The node is scheduled by week or month, and is waiting for the scheduling time to arrive. The scheduling system regularly runs the node but sets the node status to Succeeded when the scheduling system starts to run the node.
// CONDITION_UNCHOOSE(6): The node is not selected by its ancestor branch node and is run as a dry-run node.
// REALTIME_DEPRECATED(7): The node has instances that are generated in real time but are deprecated. The scheduling system sets the node status to Succeeded.
System.out.println("taskType: "+dataParam.getString("taskType"));
// The time when the node instance was modified.
System.out.println("modifyTime: "+dataParam.getString("modifyTime"));
// The time when the node instance was created.
System.out.println("createTime: "+dataParam.getString("createTime"));
// The ID of the workspace. You can call the ListProjects operation to query the ID.
System.out.println("appId: "+dataParam.getString("appId"));
// The ID of the tenant that manages the workspace to which the auto triggered node instance belongs.
System.out.println("tenantId: "+dataParam.getString("tenantId"));
// The operation code of the auto triggered node instance. You can ignore the field value.
System.out.println("opCode: "+dataParam.getString("opCode"));
// The ID of the workflow. For an auto triggered node instance, the field value is 1. For a manually triggered workflow or an auto triggered node instance of the internal workflow type, the field value is the actual workflow ID.
System.out.println("flowId: "+dataParam.getString("flowId"));
// The ID of the node for which the auto triggered node instance was generated.
System.out.println("nodeId:"+dataParam.getString("nodeId"));
// The time when the auto triggered node instance started to wait for resources.
System.out.println("beginWaitResTime: "+dataParam.getString("beginWaitResTime"));
// The ID of the auto triggered node instance.
System.out.println("taskId: "+dataParam.getString("taskId"));
// The status of the node. Valid values:
// 0: The node is not running.
// 2: The node is waiting for the scheduling time to arrive. The scheduling time is specified by the dueTime or cycleTime parameter.
// 3: The node is waiting for resources.
// 4: The node is running.
// 7: The tables that are specified in the node are issued to Data Quality and data in the tables is checked based on monitoring rules in Data Quality.
// 8: Branch conditions are being checked.
// 5: The node failed to be run.
// 6: The node is successfully run.
System.out.println("status: "+dataParam.getString("status"));
}else{
System.out.println("Failed to filter out other types of events. Check the parameter configurations.");
}
}
}
Déploiement d'exemple de projet
-
Préparez l'environnement et le projet.
Configuration requise pour l'environnement : Java 8 ou version ultérieure et Maven. Maven est un outil d'automatisation de build pour Java.
Lien de téléchargement du fichier projet : event-demo-instance-status.zip.
-
Après avoir téléchargé le fichier projet, accédez au répertoire racine du projet et exécutez la commande suivante :
mvn clean package -Dmaven.test.skip=true spring-boot:repackage -
Une fois que vous avez obtenu le package JAR pouvant être installé directement, exécutez la commande suivante :
java -jar target/event-demo-instance-status-1.0.jarLa figure suivante montre le projet démarré avec succès.

Saisissez
http://localhost:8080/indexdans la barre d'adresse d'un navigateur et appuyez sur Entrée. Si"hello world!"est renvoyé, l'extension est déployée avec succès. Vous pouvez vous abonner aux messages d'événement une fois que les connexions réseau sont établies entre DataWorks et votre extension, ainsi qu'entre EventBridge et votre extension.
Activer et configurer l'abonnement aux messages d'événement (OpenEvent)
Activer l'abonnement aux messages
Dans la console EventBridge, créez un bus personnalisé. Vous n'avez pas besoin de configurer les paramètres dans les étapes Event Source, Event Rule et Event Target.

-
Dans la console EventBridge, créez une règle d'événement pour le bus personnalisé.
Dans cet exemple, le bus personnalisé est configuré pour recevoir les messages relatifs aux événements de changement d'état d'une instance de nœud planifié. Le contenu suivant fournit un exemple de configuration d'une démo et des paramètres principaux pour la règle d'événement :
-
Configurez les paramètres à l'étape Configure Event Pattern.

{ "source": [ "acs.dataworks" ], "type": [ "dataworks:InstanceStatusChanges:InstanceStatusChanges" ] }source : l'identifiant du service dans lequel un événement se produit. Définissez ce paramètre sur acs.dataworks.
type : le type de l'événement qui se produit dans le service. Définissez ce paramètre sur dataworks:InstanceStatusChanges:InstanceStatusChanges. Modifiez les valeurs des paramètres source et type dans la section Event Pattern Debugging puis cliquez sur Test. Si le test réussit, cliquez sur Next Step.

À l'étape Configure Targets, définissez Service Type sur HTTPS et saisissez une URL valide. Utilisez les paramètres par défaut pour les autres paramètres.

-
Accédez à la page Open Platform dans la console DataWorks. Dans l'arborescence de navigation de gauche, cliquez sur OpenEvent. Sur la page qui s'affiche, ajoutez un canal de distribution d'événements.

Écrire le code
package com.aliyun.dataworks.demo;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.aliyun.dataworks.config.Constants;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
/**
* @author dataworks demo
*/
@RestController
@RequestMapping("/event")
public class ExtensionsController {
/**
* Receive event messages that are sent from EventBridge.
* @param jsonParam
*/
@PostMapping("/consumer")
public void consumerEventBridge(@RequestBody String jsonParam){
JSONObject jsonObj = JSON.parseObject(jsonParam);
String eventCode = jsonObj.getString(Constants.EVENT_CODE_FILED);
if(Constants.INSTANCE_STATUS_EVENT_CODE.equals(eventCode)){
JSONObject dataParam = JSON.parseObject(jsonObj.getString("data"));
// The time when the auto triggered node instance started to wait for the scheduling time.
System.out.println("beginWaitTimeTime: "+ dataParam.getString("beginWaitTimeTime"));
//DagId
System.out.println("dagId: "+ dataParam.getString("dagId"));
// The type of the directed acyclic graph (DAG). Valid values:
// 0: for auto triggered nodes
// 1: for manually triggered nodes
// 2: for smoke testing
// 3: for nodes for which you backfill data
// 4: for manually triggered workflows
// 5: for temporary workflows
System.out.println("dagType: "+dataParam.getString("dagType"));
// The type of the node. Valid values:
// NORMAL(0): The node is an auto triggered node. The scheduling system regularly runs the node.
// MANUAL(1): The node is a manually triggered node. The scheduling system does not regularly run the node.
// PAUSE(2): The node is a frozen node. The scheduling system regularly runs the node but sets the node status to Failed when the scheduling system starts to run the node.
// SKIP(3): The node is a dry-run node. The scheduling system regularly runs the node but sets the node status to Succeeded when the scheduling system starts to run the node.
// SKIP_UNCHOOSE(4): The node is an unselected node in a temporary workflow. This type of node exists only in temporary workflows. The scheduling system sets the node status to Succeeded when the scheduling system starts to run the node.
// SKIP_CYCLE(5): The node is scheduled by week or month, and is waiting for the scheduling time to arrive. The scheduling system regularly runs the node but sets the node status to Succeeded when the scheduling system starts to run the node.
// CONDITION_UNCHOOSE(6): The node is not selected by its ancestor branch node and is run as a dry-run node.
// REALTIME_DEPRECATED(7): The node has instances that are generated in real time but are deprecated. The scheduling system sets the node status to Succeeded.
System.out.println("taskType: "+dataParam.getString("taskType"));
// The time when the node instance was modified.
System.out.println("modifyTime: "+dataParam.getString("modifyTime"));
// The time when the node instance was created.
System.out.println("createTime: "+dataParam.getString("createTime"));
// The ID of the workspace. You can call the ListProjects operation to query the ID.
System.out.println("appId: "+dataParam.getString("appId"));
// The ID of the tenant that manages the workspace to which the auto triggered node instance belongs.
System.out.println("tenantId: "+dataParam.getString("tenantId"));
// The operation code of the auto triggered node instance. You can ignore the field value.
System.out.println("opCode: "+dataParam.getString("opCode"));
// The ID of the workflow. For an auto triggered node instance, the field value is 1. For a manually triggered workflow or an auto triggered node instance of the internal workflow type, the field value is the actual workflow ID.
System.out.println("flowId: "+dataParam.getString("flowId"));
// The ID of the node for which the auto triggered node instance was generated.
System.out.println("nodeId:"+dataParam.getString("nodeId"));
// The time when the auto triggered node instance started to wait for resources.
System.out.println("beginWaitResTime: "+dataParam.getString("beginWaitResTime"));
// The ID of the auto triggered node instance.
System.out.println("taskId: "+dataParam.getString("taskId"));
// The status of the node. Valid values:
// 0: The node is not running.
// 2: The node is waiting for the scheduling time to arrive. The scheduling time is specified by the dueTime or cycleTime parameter.
// 3: The node is waiting for resources.
// 4: The node is running.
// 7: The tables that are specified in the node are issued to Data Quality and data in the tables is checked based on monitoring rules in Data Quality.
// 8: Branch conditions are being checked.
// 5: The node failed to be run.
// 6: The node is successfully run.
System.out.println("status: "+dataParam.getString("status"));
}else{
System.out.println("Failed to filter out other types of events. Check the parameter configurations.");
}
}
}
Déployer et exécuter le code sur votre machine locale
-
Téléchargez le fichier du projet de démonstration :
Configuration requise pour l'environnement : Java 8 ou version ultérieure et Maven. Maven est un outil d'automatisation de build pour Java.
Lien de téléchargement du fichier projet : event-demo-instance-status.zip.
-
Méthodes de déploiement
Déploiement local : Empaquetez le projet sous forme de fichier JAR, puis exécutez la commande
java -jar yourapp.jarsur un serveur local disposant de Java 8 et Maven installés pour démarrer le service.Déploiement cloud : Empaquetez le projet sous forme de fichier JAR, puis téléchargez-le dans l'environnement d'exécution approprié, tel qu'un conteneur Docker ou un serveur cloud, pour le déploiement.
RemarqueLe service déployé doit être accessible par EventBridge via Internet.
-
Après avoir téléchargé le fichier projet, accédez au répertoire racine du projet et exécutez la commande suivante :
mvn clean package -Dmaven.test.skip=true spring-boot:repackageUne fois que vous avez obtenu le package JAR pouvant être installé directement, exécutez la commande suivante :
java -jar target/event-demo-instance-status-1.0.jarLa figure suivante montre le projet démarré avec succès.
Saisissez http://localhost:8080/indexdans la barre d'adresse d'un navigateur et appuyez sur Entrée. Si"hello world!"est renvoyé, l'extension est déployée avec succès. Vous pouvez vous abonner aux messages d'événement une fois que les connexions réseau sont établies entre DataWorks et votre extension, ainsi qu'entre EventBridge et votre extension.