Cette rubrique explique comment créer un connecteur de puits OSS pour diffuser les données d'une rubrique source dans ApsaraMQ for Kafka vers des objets dans Object Storage Service (OSS).
Prérequis
Pour plus d'informations, consultez la section Prérequis.
Notes d'utilisation
Le connecteur utilise l'heure de traitement de l'événement pour le partitionnement temporel, et non l'heure de création de l'événement. Pour le partitionnement temporel, les données proches d'une limite temporelle peuvent être acheminées vers le sous-répertoire de la partition temporelle suivante.
Gestion des données erronées : si vous configurez une expression JsonPath pour une partition personnalisée ou pour le contenu du fichier, mais qu'un message ne correspond pas à l'expression, le connecteur achemine ces données vers le chemin
invalidRuleData/du bucket en fonction de la politique de regroupement. Si vous constatez la présence de ce répertoire dans votre bucket, vérifiez l'exactitude de votre expression JsonPath et assurez-vous que les consommateurs ne perdent aucune donnée.La latence de bout en bout peut varier de quelques secondes à plusieurs minutes.
Si l'expression JsonPath configurée pour une partition personnalisée ou le contenu du fichier doit extraire des données du message source Kafka, vous devez encoder le contenu du message au format JSON côté source Kafka.
Le connecteur ajoute les données amont à OSS en temps réel. Par conséquent, le dernier fichier d'un chemin de partition est généralement en cours d'écriture et n'est pas dans un état final. Consommez ces données avec prudence.
Facturation
Les tâches du connecteur s'exécutent sur la plateforme Alibaba Cloud Function Compute. Vous payez les ressources de calcul consommées pour le traitement et la transmission des données selon le prix unitaire de Function Compute. Pour plus d'informations, consultez la section Vue d'ensemble de la facturation.
Étape 1 : Créer un bucket de destination
Créez un bucket dans la console Object Storage Service. Pour obtenir des instructions détaillées, consultez la section Créer des buckets.
Cette rubrique utilise un bucket nommé oss-sink-connector-bucket à titre d'exemple.
Étape 2 : Créer et démarrer le connecteur de puits OSS
Connectez-vous à la console ApsaraMQ for Kafka. Sur la page Overview, sélectionnez une région dans la section Resource Distribution.
Dans le volet de navigation de gauche, choisissez .
-
Sur la page Tasks, cliquez sur Create task.
-
Create task
-
À l'étape Source, définissez Data Provider sur ApsaraMQ for Kafka, configurez les paramètres suivants, puis cliquez sur Next.
Parameter
Description
Example
Region
Région où se trouve l'instance ApsaraMQ for Kafka source.
China (Beijing)
ApsaraMQ for Kafka Instance
Instance ApsaraMQ for Kafka source qui produit les messages.
alikafka_post-cn-jte3****
Topic
Sélectionnez la rubrique pour la production de messages ApsaraMQ for Kafka.
topic
Group ID
ID du groupe de consommateurs pour l'instance ApsaraMQ for Kafka source.
-
Quick Create : crée un nouvel ID de groupe.
-
Use Existing : sélectionnez un ID de groupe existant.
GID_http_1
Consumer Offset
Position à partir de laquelle commencer la consommation des messages.
Latest
Network Configuration
Type de réseau utilisé pour acheminer les messages.
Classic Network
VPC
ID du VPC. Ce paramètre est requis uniquement si Network Configuration est défini sur VPC.
vpc-bp17fapfdj0dwzjkd****
vSwitches
ID du vSwitch. Ce paramètre est requis uniquement si Network Configuration est défini sur VPC.
vsw-bp1gbjhj53hdjdkg****
Security Group
Groupe de sécurité. Ce paramètre est requis uniquement si Network Configuration est défini sur VPC.
alikafka_pre-cn-7mz2****
Messages
Nombre maximal de messages à envoyer par invocation de fonction. Une requête est envoyée uniquement lorsque le nombre de messages accumulés atteint cette valeur. Valeurs valides : 1 à 10 000.
100
Interval (Unit: Seconds)
Intervalle entre les invocations de fonction. Le système regroupe et envoie les messages à Function Compute à cet intervalle. Valeurs valides : 0 à 15 secondes. Une valeur de 0 indique que les messages sont livrés immédiatement.
3
-
À l'étape Filtering, définissez le Pattern Content pour filtrer les événements. Pour plus d'informations, consultez la section event pattern.
À l'étape Transformation, configurez la transformation des données pour effectuer des opérations telles que le fractionnement, le mappage, l'enrichissement et le routage dynamique. Pour plus d'informations, consultez la section Use Function Compute to clean message data.
-
À l'étape Sink, définissez Service Type sur OSS, et configurez les paramètres suivants.
Parameter
Description
Example
OSS bucket
Bucket OSS de destination.
Important-
Assurez-vous que le bucket spécifié existe et n'est pas supprimé pendant l'exécution de la tâche.
-
Vous pouvez livrer des données uniquement aux buckets des classes de stockage Standard ou Infrequent Access (IA). Vous ne pouvez pas livrer de données aux buckets de la classe de stockage Archive.
-
Après la création d'une tâche de connecteur de puits OSS, la plateforme crée un répertoire système
.tmp/à la racine du bucket OSS. Ne supprimez pas ce répertoire et n'utilisez pas les objets qu'il contient.
oss-sink-connector-bucket
Storage Path
Une clé d'objet OSS se compose d'un chemin et d'un nom. Par exemple, si la clé d'objet est
a/b/c/a.txt, le chemin esta/b/c/et le nom esta.txt. Vous pouvez personnaliser le chemin pour le partitionnement. Le nom est généré automatiquement par le connecteur avec le format :{unix_timestamp_in_milliseconds}_{8_character_random_string}. Exemple :1705576353794_elJmxu3v.-
Si ce paramètre n'est pas configuré ou est défini sur
/, les données ne sont pas partitionnées et sont enregistrées dans le répertoire racine du bucket. -
Les variables temporelles sont prises en charge :
{yyyy},{MM},{dd}et{HH}représentent respectivement l'année, le mois, le jour et l'heure. Ces variables sont sensibles à la casse. -
Les expressions JsonPath sont prises en charge pour personnaliser les paramètres de chemin, tels que
{$.data.topic}et{$.data.partition}. Les expressions doivent être valides. En raison des restrictions de chemin OSS, nous recommandons d'extraire les valeurs sous forme d'intou destring. Les valeurs extraites doivent être des caractères UTF-8 standard et ne doivent pas contenir d'espaces, de.., d'emojis, de/ou de\afin d'éviter les erreurs d'écriture des données. -
Les constantes sont prises en charge.
RemarqueLes partitions vous permettent de regrouper les données logiquement et d'éviter les problèmes causés par un nombre excessif de petits fichiers dans un seul chemin.
Le débit du connecteur évolue avec le nombre de partitions. Trop peu de partitions, ou l'absence de partitions, peuvent réduire le débit du connecteur et provoquer des retards en amont. Trop de partitions peuvent entraîner une fragmentation des données, une augmentation des opérations d'écriture et de nombreux petits fichiers. Par conséquent, votre stratégie de partitionnement est cruciale. Tenez compte des recommandations suivantes :
-
Source Kafka : partitionnez à la fois par heure et par ID de partition Kafka. Si les performances sont insuffisantes, augmentez le nombre de partitions Kafka pour améliorer le débit du connecteur. Exemple :
prefix/{yyyy}/{MM}/{dd}/{HH}/{$.data.partition}/ -
Regroupement basé sur l'activité : partitionnez par un champ métier dans les données. Le débit dépend du nombre de valeurs distinctes pour ce champ. Exemple :
prefixV2/{$.data.body.field}/
Nous vous recommandons d'utiliser différents préfixes constants pour différentes tâches. Cela évite que plusieurs tâches partagent des partitions et écrivent dans le même répertoire, ce qui pourrait mélanger les données et compliquer la gestion.
-
alikafka_post-cn-9dhsaassdd****/guide-oss-sink-topic/YYYY/MM/dd/HH
Time Zone
Le fuseau horaire par défaut est UTC+8:00. Ce paramètre s'applique uniquement au partitionnement temporel.
UTC+8:00
Batch aggregation object size
Taille cible pour le regroupement des données dans un seul objet. Valeurs valides : 1 à 1 024 Mo.
Remarque-
Le connecteur écrit les données dans un objet OSS par lots. La taille de chaque lot se situe dans la plage (0 Mo, 16 Mo]. Par conséquent, la taille finale de l'objet peut être légèrement supérieure à la valeur configurée, jusqu'à 16 Mo.
-
Dans les scénarios à fort trafic, nous vous recommandons de définir la taille de l'objet de regroupement à plusieurs centaines de mégaoctets (par exemple, 128 Mo ou 512 Mo) et la fenêtre de temps de regroupement au niveau de l'heure (par exemple, 60 min ou 120 min).
5
Batch aggregation time window
Fenêtre de temps pour le regroupement des données. Valeurs valides : 1 à 1 440 minutes.
1
File Compression
-
None : l'objet OSS généré n'a pas d'extension de fichier.
-
GZIP : l'objet généré a une extension .gz.
-
Snappy : l'objet généré a une extension .snappy.
-
Zstd : l'objet généré a une extension .zstd.
Lorsque la compression est activée, le connecteur regroupe les données en fonction de la taille avant compression. Par conséquent, la taille de l'objet affichée dans OSS est inférieure à la taille de lot configurée, et la taille décompressée est proche de la taille de lot.
None
File Content
-
Full Event : le connecteur encapsule le message original dans l'enveloppe CloudEvents. Les données résultantes incluent les métadonnées CloudEvents. Dans l'exemple suivant, le champ
datacontient le message original, et les autres champs sont des métadonnées de l'enveloppe CloudEvents.{ "specversion": "1.0", "id": "8e215af8-ca18-4249-8645-f96c1026****", "source": "acs:alikafka", "type": "alikafka:Topic:Message", "subject": "acs:alikafka:alikafka_pre-cn-i7m2msb9****:topic:****", "datacontenttype": "application/json; charset=utf-8", "time": "2022-06-23T02:49:51.589Z", "aliyunaccountid": "182572506381****", "data": { "topic": "****", "partition": 7, "offset": 25, "timestamp": 1655952591589, "headers": { "headers": [], "isReadOnly": false }, "key": "keytest", "value": "hello kafka msg" } } -
Partial Event : livre uniquement une partie des données, extraite par une expression JsonPath. Par exemple, si vous configurez
$.data, seule la valeur du champdataest livrée à OSS.
Pour exclure les métadonnées CloudEvents, sélectionnez Partial Event et utilisez l'expression
$.data. Cela permet de livrer uniquement le message source original à OSS, ce qui réduit les coûts de stockage et améliore l'efficacité de la transmission.Partial Event
$.data -
-
-
Task properties
Configurez la politique de nouvelle tentative pour les livraisons d'événements ayant échoué et la méthode de gestion des erreurs. Pour plus d'informations, consultez la section Retry policies and dead-letter queues.
-
Revenez à la page Tasks, recherchez votre tâche et cliquez sur Enable dans la colonne Actions.
-
Dans la boîte de dialogue Note, lisez l'invite et cliquez sur OK.
Le démarrage de la tâche prend entre 30 et 60 secondes. Vous pouvez consulter la progression du démarrage dans la colonne Status de la page Tasks.
Étape 3 : Vérifier le connecteur de puits OSS
Sur la page Message Outflow, dans la ligne correspondant à votre tâche de connecteur de puits OSS, cliquez sur le nom de la rubrique source dans la colonne Event Sources.
Sur la page des détails de la rubrique, cliquez sur Send Test Message.
-
Dans le panneau Quickly experience message sending and consumption, configurez le message comme suit et cliquez sur OK.
Définissez Message key sur
oss-sink-k2, définissez Message content suross-sink-v2, et pour Send to specified partition, sélectionnez No. Sur la page Message Outflow, dans la ligne correspondant à votre tâche de connecteur de puits OSS, cliquez sur le nom du bucket de destination dans la colonne Event Targets.
-
Sur la page Bucket, dans la barre de navigation de gauche, sélectionnez .
Répertoire
tmp/: répertoire système requis par le connecteur. Ne supprimez pas et n'utilisez pas les objets de ce répertoire.Répertoires de données : les sous-répertoires sont générés en fonction de la règle de chemin de partition configurée pour la tâche. Les fichiers de données sont téléchargés dans le répertoire le plus profond.
La structure de chemin du répertoire de données est similaire à
alikafka_p[TopicName]/[PartitionID]/yyyy/MM/dd/HH/. Dans le répertoire le plus profond, vous trouverez un fichier de métadonnées (tel que.oss_meta_file...) et des fichiers de données de partition (tels quepartition_3_of...). La présence de ces fichiers confirme que le connecteur de puits OSS a bien écrit les messages dans OSS. Dans la colonne Actions de l'objet, choisissez .
-
Ouvrez le fichier téléchargé pour afficher le contenu du message.
{"topic":"guide-oss-sink-topic","partition":0,"offset":0,"timestamp":1681378474218,"headers":{"headers":[],"isReadOnly":false},"key":"oss-sink-k2","value":"oss-sink-v2"} {"topic":"guide-oss-sink-topic","partition":0,"offset":1,"timestamp":1681378491498,"headers":{"headers":[],"isReadOnly":false},"key":"oss-sink-k2","value":"oss-sink-v2"} {"topic":"guide-oss-sink-topic","partition":0,"offset":2,"timestamp":1681378492515,"headers":{"headers":[],"isReadOnly":false},"key":"oss-sink-k2","value":"oss-sink-v2"}Les messages sont séparés par des sauts de ligne.