Realtime Compute for Apache Flink prend en charge la lecture et l'écriture dans le service Object Storage Service (OSS) à l'aide du connecteur de système de fichiers. OSS offre une durabilité des données de 99,9999999999 % (douze neufs) et une disponibilité de 99,995 %, ce qui en fait un stockage fiable pour les pipelines Flink à grande échelle.
| Catégorie | Détails |
|---|---|
| Types de table pris en charge | Tables source et puits |
| Modes d'exécution | Traitement par lots et en continu |
| Formats de données | ORC, Parquet, Avro, CSV, JSON et brut |
| Métriques de surveillance spécifiques | Aucune |
| Types d'API | API DataStream et SQL |
| Mise à jour ou suppression de données dans les tables puits | Insertion uniquement. La mise à jour et la suppression ne sont pas prises en charge. |
Limites
Général :
Seules les versions Ververica Runtime (VVR) 11 et ultérieures prennent en charge la lecture des fichiers compressés (GZIP, BZIP2, XZ, DEFLATE) depuis OSS. VVR 8 ne peut pas traiter correctement les fichiers compressés.
Les versions de VVR antérieures à la 8.0.6 prennent uniquement en charge les buckets OSS situés dans le même compte. Pour accéder aux buckets entre différents comptes, utilisez VVR 8.0.6 ou une version ultérieure et configurez l'authentification du bucket. Pour plus de détails, consultez la section Configurer l'authentification du bucket.
La lecture incrémentielle des nouvelles partitions n'est pas prise en charge.
L'accès cross-région à OSS n'est pas pris en charge. Le bucket OSS doit se trouver dans la même région que l'espace de travail Flink. Le connecteur OSS est basé sur l'interface Filesystem, qui ne prend en charge qu'un seul endpoint global. L'accès à un bucket OSS situé dans une autre région provoque une erreur de correspondance d'endpoint.
Tables puits uniquement :
Les formats de ligne Avro, CSV, JSON et brut ne sont pas pris en charge lors de l'écriture dans OSS. Consultez FLINK-30635 pour plus de détails.
Syntaxe
CREATE TABLE OssTable (
column_name1 INT,
column_name2 STRING,
...
datetime STRING,
`hour` STRING
) PARTITIONED BY (datetime, `hour`) WITH (
'connector' = 'filesystem', -- required: must be 'filesystem'
'path' = 'oss://<bucket>/path', -- required: URI of the OSS path
'format' = '...', -- required: orc, parquet, avro, csv, json, or raw
'partition.default-name' = '...', -- optional: partition name when partition field is NULL or empty
'source.monitor-interval' = '...', -- optional (source only): interval to scan for new files
'auto-compaction' = '...' -- optional (sink only): enable automatic compaction after each checkpoint
);
Colonnes de métadonnées
Les tables sources prennent en charge les colonnes de métadonnées qui exposent les informations au niveau du fichier pour chaque ligne. Définissez une colonne de métadonnées dans votre DDL en ajoutant METADATA après le type de données :
CREATE TABLE MyUserTableWithFilepath (
column_name1 INT,
column_name2 STRING,
`file.path` STRING NOT NULL METADATA
) WITH (
'connector' = 'filesystem',
'path' = 'oss://<bucket>/path',
'format' = 'json'
)
Les colonnes de métadonnées suivantes sont disponibles :
| Clé | Type de données | Description |
|---|---|---|
file.path |
STRING NOT NULL | Chemin complet du fichier contenant la ligne. |
file.name |
STRING NOT NULL | Nom du fichier (dernier élément du chemin). |
file.size |
BIGINT NOT NULL | Taille du fichier, en octets. |
file.modification-time |
TIMESTAMP_LTZ(3) NOT NULL | Dernière heure de modification du fichier. |
Paramètres WITH
Paramètres généraux
| Paramètre | Obligatoire | Par défaut | Description |
|---|---|---|---|
connector |
Oui | — | Doit être filesystem. |
path |
Oui | — | Chemin OSS au format URI, tel que oss://my_bucket/my_path. Pour VVR 8.0.6 et versions ultérieures, l'authentification du bucket est requise après la définition de ce paramètre. Consultez Configurer l'authentification du bucket. |
format |
Oui | — | Format de fichier : csv, json, avro, parquet, orc ou raw. |
Paramètres de la table source
| Paramètre | Obligatoire | Par défaut | Description |
|---|---|---|---|
source.monitor-interval |
Non | — | Intervalle d'analyse pour détecter les nouveaux fichiers. Doit être supérieur à 0. Si ce paramètre n'est pas défini, le chemin est analysé une seule fois et la source est bornée. Chaque fichier est identifié par son chemin et traité exactement une fois. Les chemins des fichiers traités sont stockés dans l'état et conservés lors des points de contrôle et des points de sauvegarde. Un intervalle plus court accélère la découverte des fichiers mais augmente la fréquence d'analyse. |
Paramètres de la table puits
| Paramètre | Obligatoire | Par défaut | Description |
|---|---|---|---|
partition.default-name |
Non | _DEFAULT_PARTITION__ |
Nom de partition utilisé lorsqu'un champ de partition est NULL ou une chaîne vide. |
sink.rolling-policy.file-size |
Non | 128 Mo | Taille maximale du fichier partiel avant roulement. Chaque sous-tâche de puits crée au moins un fichier partiel par partition. Consultez Comportement de la politique de roulement pour savoir comment cela interagit avec les formats de fichiers. |
sink.rolling-policy.rollover-interval |
Non | 30 min | Durée maximale pendant laquelle un fichier partiel peut rester ouvert avant roulement. La fréquence de vérification est contrôlée par sink.rolling-policy.check-interval. |
sink.rolling-policy.check-interval |
Non | 1 min | Fréquence à laquelle il faut vérifier si un fichier partiel doit être roulé en fonction de sink.rolling-policy.rollover-interval. |
auto-compaction |
Non | false | Indique s'il faut activer la compaction automatique. Les données sont d'abord écrites dans des fichiers temporaires. Après chaque point de contrôle, les fichiers temporaires de ce point de contrôle sont fusionnés en fichiers plus volumineux. Les fichiers temporaires ne sont pas visibles avant la fusion. Lorsque cette option est activée : seuls les fichiers d'un point de contrôle sont fusionnés (au moins un fichier par point de contrôle) ; la latence de visibilité des données équivaut à checkpoint interval + compaction duration ; les longues exécutions de compaction peuvent provoquer une contre-pression et retarder les points de contrôle. |
compaction.file-size |
Non | 128 Mo | Taille cible du fichier pour la sortie compactée. La valeur par défaut est identique à celle de sink.rolling-policy.file-size. |
sink.partition-commit.trigger |
Non | process-time |
Moment auquel une partition est validée. Consultez Déclencheurs de validation de partition. |
sink.partition-commit.delay |
Non | 0s |
Délai minimal avant la validation d'une partition. Définissez sur 1 d pour les partitions quotidiennes et 1 h pour les partitions horaires. |
sink.partition-commit.watermark-time-zone |
Non | UTC |
Fuseau horaire utilisé pour analyser un watermark LONG en TIMESTAMP afin de comparer la validation de la partition. S'applique uniquement lorsque sink.partition-commit.trigger est défini sur partition-time. Utilisez le fuseau horaire de la session lorsque le watermark est défini sur une colonne TIMESTAMP_LTZ (par exemple, Asia/Shanghai). Accepte les noms complets de fuseaux horaires (tels que America/Los_Angeles) ou les décalages personnalisés (tels que GMT-08:00). Si ce paramètre n'est pas configuré correctement, les validations de partition peuvent être retardées de plusieurs heures. |
partition.time-extractor.kind |
Non | default |
Méthode d'extraction de l'heure à partir des champs de partition. default : configurez un modèle ou un formateur d'horodatage. custom : spécifiez une classe d'extracteur. |
partition.time-extractor.class |
Non | — | Classe qui implémente l'interface PartitionTimeExtractor. Obligatoire lorsque partition.time-extractor.kind est défini sur custom. |
partition.time-extractor.timestamp-pattern |
Non | — | Modèle permettant de construire un horodatage à partir des champs de partition. Par défaut, le premier champ est extrait en utilisant yyyy-MM-dd hh:mm:ss. Exemples : $dt (champ unique), $year-$month-$day $hour:00:00 (champs multiples), $dt $hour:00:00 (deux champs). |
partition.time-extractor.timestamp-formatter |
Non | yyyy-MM-dd HH:mm:ss |
Formateur permettant de convertir la chaîne d'horodatage (telle qu'exprimée par partition.time-extractor.timestamp-pattern) en horodatage. Par exemple, si partition.time-extractor.timestamp-pattern est $year$month$day, définissez ce paramètre sur yyyyMMdd. Compatible avec DateTimeFormatter de Java. |
sink.partition-commit.policy.kind |
Non | — | Méthode utilisée pour informer les consommateurs en aval qu'une partition est prête. success-file : écrit un fichier _SUCCESS dans le répertoire de la partition. custom : utilise une classe implémentant PartitionCommitPolicy. Plusieurs politiques peuvent être combinées. |
sink.partition-commit.policy.class |
Non | — | Classe qui implémente PartitionCommitPolicy. Obligatoire lorsque sink.partition-commit.policy.kind est défini sur custom. |
sink.partition-commit.success-file.name |
Non | _SUCCESS |
Nom du fichier de succès écrit par la politique de validation success-file. |
sink.parallelism |
Non | — | Parallélisme de l'opérateur d'écriture de fichiers. La valeur par défaut est le parallélisme de l'opérateur en amont. Doit être supérieur à 0. Lorsque auto-compaction est activé, l'opérateur de compaction utilise également ce parallélisme. |
Comportement de la politique de roulement
Le comportement de roulement varie selon le format de fichier :
Pour les formats columnaires (Parquet, ORC, Avro), un fichier partiel est toujours roulé au moment du point de contrôle, même si les critères de la politique de roulement ne sont pas remplis. La taille du fichier et l'intervalle de roulement s'appliquent comme déclencheurs supplémentaires entre les points de contrôle.
Pour les formats de ligne (CSV, JSON, brut), un fichier partiel est roulé uniquement lorsque les critères de la politique de roulement (sink.rolling-policy.file-sizeousink.rolling-policy.rollover-interval) sont remplis. Si vous avez besoin d'une faible latence de visibilité des fichiers, ajustezsink.rolling-policy.rollover-intervalconjointement avec votre intervalle de point de contrôle.
Les formats de ligne ne sont pas pris en charge pour les tables puits OSS en raison de FLINK-30635 . Le comportement ci-dessus s'applique si la prise en charge des formats de ligne est ajoutée dans une version future.
Déclencheurs de validation de partition
Deux types de déclencheurs sont disponibles pour sink.partition-commit.trigger :
process-time(par défaut) : Valide une partition lorsque l'heure système actuelle dépasse l'heure de création de la partition plussink.partition-commit.delay. Ne nécessite ni watermark ni extracteur de temps de partition. Plus général mais moins précis : les retards ou les pannes de données peuvent entraîner des validations prématurées.partition-time: Valide une partition lorsque le watermark dépasse l'heure de création de la partition plussink.partition-commit.delay. Nécessite la génération de watermarks et des partitions basées sur le temps (horaire, quotidien, etc.).
Configurer l'authentification du bucket
Seules les versions VVR 8.0.6 et ultérieures prennent en charge l'authentification du bucket.
Après avoir défini le paramètre path, configurez l'authentification du bucket afin que Flink puisse lire et écrire dans le chemin OSS spécifié. Ajoutez les éléments suivants à la section Additional Configurations de l'onglet Parameters de la page Deployment Details dans la console de développement Realtime Compute :
fs.oss.bucket.<bucketName>.accessKeyId: <your-access-key-id>
fs.oss.bucket.<bucketName>.accessKeySecret: <your-access-key-secret>
Remplacez <bucketName> par le nom du bucket utilisé dans le paramètre path.
| Élément de configuration | Description |
|---|---|
fs.oss.bucket.<bucketName>.accessKeyId |
AccessKey ID pour le bucket. Utilisez un AccessKey existant ou créez-en un. Consultez Créer un AccessKey. |
fs.oss.bucket.<bucketName>.accessKeySecret |
AccessKey Secret pour le bucket. |
L'AccessKey Secret n'est affiché qu'une seule fois lors de sa création. Stockez-le en toute sécurité.
Écrire dans OSS-HDFS
Ajoutez la configuration suivante à la section Other Configuration de l'onglet Parameters de la page Deployment Details dans la console de développement Realtime Compute :
fs.oss.jindo.buckets: <bucket-names>
fs.oss.jindo.accessKeyId: <your-access-key-id>
fs.oss.jindo.accessKeySecret: <your-access-key-secret>
| Élément de configuration | Description |
|---|---|
fs.oss.jindo.buckets |
Noms des buckets OSS-HDFS, séparés par des points-virgules. Lorsque Flink écrit dans un chemin OSS, si le bucket correspondant figure ici, les données sont écrites dans le service OSS-HDFS. |
fs.oss.jindo.accessKeyId |
AccessKey ID. Consultez Créer un AccessKey. |
fs.oss.jindo.accessKeySecret |
AccessKey Secret. |
L'AccessKey Secret n'est affiché qu'une seule fois lors de sa création. Stockez-le en toute sécurité.
Configurez l'endpoint OSS-HDFS en utilisant l'une des méthodes suivantes :
Configuration des paramètres
Ajoutez l'endpoint à Additional Configurations :
fs.oss.jindo.endpoint: <oss-hdfs-endpoint>
Configuration du chemin
Intégrez l'endpoint directement dans le chemin OSS :
oss://<bucket-name>.<oss-hdfs-endpoint>/<directory>
Lorsque vous utilisez cette méthode, fs.oss.jindo.buckets doit inclure <bucket-name>.<oss-hdfs-endpoint>.
Par exemple, si le nom du bucket est jindo-test et que l'endpoint est cn-beijing.oss-dls.aliyuncs.com :
# OSS path
oss://jindo-test.cn-beijing.oss-dls.aliyuncs.com/<directory>
# Additional Configurations
fs.oss.jindo.buckets: jindo-test,jindo-test.cn-beijing.oss-dls.aliyuncs.com
Écriture dans un système de fichiers distribué Hadoop (HDFS) externe
Pour les chemins utilisant le schéma hdfs://, ajoutez les éléments suivants pour spécifier ou changer le nom d'utilisateur d'accès :
containerized.taskmanager.env.HADOOP_USER_NAME: hdfs
containerized.master.env.HADOOP_USER_NAME: hdfs
Exemples
Lire depuis OSS (table source)
CREATE TEMPORARY TABLE fs_table_source (
`id` INT,
`name` VARCHAR
) WITH (
'connector' = 'filesystem',
'path' = 'oss://<bucket>/path',
'format' = 'parquet'
);
CREATE TEMPORARY TABLE blackhole_sink (
`id` INT,
`name` VARCHAR
) WITH (
'connector' = 'blackhole'
);
INSERT INTO blackhole_sink SELECT * FROM fs_table_source;
Écrire dans OSS (table puits)
Écrire dans une table partitionnée
Cet exemple diffuse les données d'une source datagen, les partitionne par date et heure, et valide les partitions à l'aide du déclencheur partition-time :
CREATE TABLE datagen_source (
user_id STRING,
order_amount DOUBLE,
ts BIGINT, -- Timestamp in milliseconds
ts_ltz AS TO_TIMESTAMP_LTZ(ts, 3),
WATERMARK FOR ts_ltz AS ts_ltz - INTERVAL '5' SECOND -- Watermark on TIMESTAMP_LTZ column
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE fs_table_sink (
user_id STRING,
order_amount DOUBLE,
dt STRING,
`hour` STRING
) PARTITIONED BY (dt, `hour`) WITH (
'connector' = 'filesystem',
'path' = 'oss://<bucket>/path',
'format' = 'parquet',
'partition.time-extractor.timestamp-pattern' = '$dt $hour:00:00',
'sink.partition-commit.delay' = '1 h',
'sink.partition-commit.trigger' = 'partition-time',
'sink.partition-commit.watermark-time-zone' = 'Asia/Shanghai',
'sink.partition-commit.policy.kind' = 'success-file'
);
INSERT INTO fs_table_sink
SELECT
user_id,
order_amount,
DATE_FORMAT(ts_ltz, 'yyyy-MM-dd'),
DATE_FORMAT(ts_ltz, 'HH')
FROM datagen_source;
Écrire dans une table non partitionnée
CREATE TABLE datagen_source (
user_id STRING,
order_amount DOUBLE
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE fs_table_sink (
user_id STRING,
order_amount DOUBLE
) WITH (
'connector' = 'filesystem',
'path' = 'oss://<bucket>/path',
'format' = 'parquet'
);
INSERT INTO fs_table_sink SELECT * FROM datagen_source;
API DataStream
Pour utiliser l'API DataStream, configurez d'abord le connecteur DataStream. Consultez Utiliser un connecteur DataStream.
L'exemple suivant utilise StreamingFileSink avec OnCheckpointRollingPolicy pour écrire dans OSS. Les fichiers partiels sont roulés à chaque point de contrôle.
String outputPath = "oss://<bucket>/path";
final StreamingFileSink<Row> sink =
StreamingFileSink.forRowFormat(
new Path(outputPath),
(Encoder<Row>) (element, stream) -> {
out.println(element.toString());
})
.withRollingPolicy(OnCheckpointRollingPolicy.build())
.build();
outputStream.addSink(sink);
Pour écrire dans OSS-HDFS, configurez également les paramètres OSS-HDFS dans Additional Configurations. Consultez Écrire dans OSS-HDFS.