Stream Load est une méthode synchrone basée sur HTTP permettant de charger directement des fichiers locaux ou des flux de données dans StarRocks. Envoyez une requête PUT : la réponse HTTP vous indique immédiatement si le chargement a réussi. Stream Load prend en charge les formats CSV et JSON, avec une limite de 10 Go par chargement unique.
Soumettre une tâche de chargement
Stream Load utilise curl pour envoyer une requête PUT au nœud FE. Vous pouvez également utiliser d'autres clients HTTP.
Syntaxe
curl --location-trusted -u <username>:<password> \
-H "Expect:100-continue" \
-XPUT <url> \
(data_desc) \
[opt_properties]
Incluez l'en-tête Expect:100-continue dans votre requête. Cela évite un transfert de données inutile si le serveur rejette la tâche avant l'envoi du corps de la requête.
URL
http://<fe_host>:<fe_http_port>/api/<database_name>/<table_name>/_stream_load
| Paramètre | Obligatoire | Description |
|---|---|---|
<fe_host> |
Oui | Adresse IP du nœud FE |
<fe_http_port> |
Oui | Port HTTP du nœud FE. Par défaut : 18030. Exécutez SHOW FRONTENDS pour récupérer l'adresse et le port, ou consultez le paramètre http_port dans l'onglet Configuration de la page Cluster Service. |
<database_name> |
Oui | Nom de la base de données contenant la table cible |
<table_name> |
Oui | Nom de la table cible |
Paramètres des données sources
Le bloc data_desc décrit le fichier source et contrôle la manière dont StarRocks mappe et filtre les données.
Paramètres communs
| Paramètre | Obligatoire | Description |
|---|---|---|
-T <file_path> |
Oui | Chemin d'accès au fichier de données source. Le nom du fichier peut inclure ou omettre l'extension. |
format |
Non | Format du fichier. Valeurs valides : CSV (par défaut), JSON |
partitions |
Non | Partitions cibles. Si omis, les données sont chargées dans toutes les partitions de la table. |
temporary_partitions |
Non | Partitions temporaires cibles |
columns |
Non | Mappage des colonnes entre le fichier source et la table StarRocks. Requis lorsque le schéma du fichier ne correspond pas au schéma de la table. Consultez les exemples de mappage de colonnes. |
Paramètres CSV
| Paramètre | Obligatoire | Description |
|---|---|---|
column_separator |
Non | Délimiteur de colonne. Par défaut : \t. Pour les caractères invisibles, utilisez la notation hexadécimale avec le préfixe \x, par exemple -H "column_separator:\x01" pour le délimiteur Hive. |
row_delimiter |
Non | Délimiteur de ligne. Par défaut : \n. Remarque
|
skip_header |
Non | Nombre de lignes d'en-tête à ignorer au début du fichier CSV. Par défaut : 0. |
where |
Non | Condition de filtrage. Seules les lignes correspondant à la condition sont chargées. Exemple : -H "where: k1 = 20180601" |
max_filter_ratio |
Non | Ratio maximal de lignes pouvant être filtrées en raison de problèmes de qualité des données. Par défaut : 0 (tolérance zéro). Les lignes filtrées par where ne sont pas comptabilisées. |
strict_mode |
Non | Contrôle le filtrage lors des conversions de type. false (par défaut) : désactivé. true : filtrage strict appliqué lors des conversions de type de colonne. |
timeout |
Non | Délai d'expiration du chargement en secondes. Par défaut : 600. Plage valide : 1–259200. |
timezone |
Non | Fuseau horaire pour le chargement. Par défaut : UTC+8. Affecte toutes les fonctions liées au fuseau horaire lors du chargement. |
exec_mem_limit |
Non | Limite de mémoire pour le chargement. Par défaut : 2 Go. |
Paramètres JSON
| Paramètre | Obligatoire | Description |
|---|---|---|
jsonpaths |
Non | Champs à charger, au format JSON. Requis uniquement pour le mode de correspondance JSON. |
strip_outer_array |
Non | Indique s'il faut supprimer le tableau JSON externe. false (par défaut) : le tableau entier est importé comme une seule valeur. true : chaque élément du tableau est importé sous forme de ligne distincte. |
json_root |
Non | Élément racine des données JSON. Requis uniquement pour le mode de correspondance JSON. Doit être une chaîne JsonPath valide. Par défaut : vide (le fichier JSON entier est importé). |
ignore_json_size |
Non | Indique s'il faut ignorer la vérification de la taille du corps JSON. Par défaut, le corps JSON ne peut pas dépasser 100 Mo. Définissez la valeur sur true pour ignorer cette vérification, mais notez que des charges volumineuses peuvent entraîner une utilisation élevée de la mémoire. |
compression / Content-Encoding |
Non | Algorithme de compression pour la transmission des données. Pris en charge : GZIP, BZIP2, LZ4_FRAME, Zstandard. |
Options de tâche de chargement
Les opt_properties sont des paramètres optionnels qui s'appliquent à l'ensemble de la tâche de chargement.
-H "label: <label_name>"
-H "where: <condition>"
-H "max_filter_ratio: <num>"
-H "timeout: <num>"
-H "strict_mode: true | false"
-H "timezone: <string>"
-H "load_mem_limit: <num>"
-H "partial_update: true | false"
-H "partial_update_mode: row | column"
-H "merge_condition: <column_name>"
| Paramètre | Obligatoire | Description |
|---|---|---|
label |
Non | Étiquette pour la tâche de chargement. Évite les importations en double : StarRocks rejette toute tâche réutilisant une étiquette issue d'une tâche réussie au cours des 30 dernières minutes. |
where |
Non | Condition de filtrage appliquée aux données transformées. Seules les lignes correspondant à la condition sont chargées. |
max_filter_ratio |
Non | Ratio maximal de lignes pouvant être filtrées en raison de problèmes de qualité des données. Par défaut : 0. Les lignes filtrées par where ne sont pas comptabilisées. |
log_rejected_record_num |
Non | Nombre maximal de lignes filtrées à journaliser par nœud BE (ou CN). Pris en charge dans StarRocks 3.1 et versions ultérieures. Valeurs valides : 0 (par défaut, aucune journalisation), -1 (journaliser tout) ou un entier positif. |
timeout |
Non | Délai d'expiration du chargement en secondes. Par défaut : 600. Plage valide : 1–259200. |
strict_mode |
Non | false (par défaut) : désactivé. true : filtrage strict appliqué. |
timezone |
Non | Fuseau horaire pour le chargement. Par défaut : UTC+8. |
load_mem_limit |
Non | Limite de mémoire pour la tâche de chargement. Par défaut : 2 Go. |
partial_update |
Non | Indique s'il faut utiliser les mises à jour partielles de colonnes. Par défaut : FALSE. |
partial_update_mode |
Non | Mode des mises à jour partielles. row (par défaut) : adapté aux mises à jour en temps réel de nombreuses colonnes par petits lots. column : adapté aux mises à jour par lot de quelques colonnes sur de nombreuses lignes ; par exemple, mettre à jour 10 colonnes sur 100 pour toutes les lignes peut améliorer les performances d'un facteur 10. |
merge_condition |
Non | Colonne utilisée comme condition de mise à jour. La mise à jour prend effet uniquement lorsque la valeur importée est supérieure ou égale à la valeur actuelle. Doit être une colonne non clé primaire. Pris en charge uniquement pour les tables à clé primaire. |
Dans SQL StarRocks, les mots clés réservés doivent être entourés d'accents graves. Consultez Keywords .
Valeur de retour
Une fois le chargement terminé, StarRocks renvoie le résultat au format JSON :
{
"TxnId": 9,
"Label": "label2",
"Status": "Success",
"Message": "OK",
"NumberTotalRows": 4,
"NumberLoadedRows": 4,
"NumberFilteredRows": 0,
"NumberUnselectedRows": 0,
"LoadBytes": 45,
"LoadTimeMs": 235,
"BeginTxnTimeMs": 101,
"StreamLoadPlanTimeMs": 102,
"ReadDataTimeMs": 0,
"WriteDataTimeMs": 11,
"CommitAndPublishTimeMs": 19
}
| Champ | Description |
|---|---|
TxnId |
ID de transaction du chargement |
Label |
Étiquette de la tâche de chargement |
Status |
Statut final du chargement. Voir le tableau ci-dessous. |
ExistingJobStatus |
Statut de la tâche existante lorsque Status est Label Already Exists. Valeurs : RUNNING ou FINISHED. |
Message |
Détails sur le statut du chargement. En cas d'échec, contient la raison de l'erreur. |
NumberTotalRows |
Nombre total de lignes lues depuis le flux de données |
NumberLoadedRows |
Lignes chargées avec succès. Valide uniquement lorsque Status est Success. |
NumberFilteredRows |
Lignes filtrées en raison de problèmes de qualité des données |
NumberUnselectedRows |
Lignes filtrées par la condition where |
LoadBytes |
Taille du fichier source |
LoadTimeMs |
Temps total de la tâche de chargement, en millisecondes |
BeginTxnTimeMs |
Temps nécessaire pour démarrer la transaction, en millisecondes |
StreamLoadPlanTimeMs |
Temps nécessaire pour générer le plan d'exécution, en millisecondes |
ReadDataTimeMs |
Temps de lecture des données, en millisecondes |
WriteDataTimeMs |
Temps d'écriture des données, en millisecondes |
ErrorURL |
URL des détails d'erreur en cas d'échec du chargement. Consultez Vérifier les détails des erreurs. |
Statut du chargement
| Statut | Signification | Action |
|---|---|---|
Success |
Données chargées et visibles | Aucune |
Publish Timeout |
Chargement soumis avec succès, mais la visibilité des données peut être retardée | Aucun nouvel essai nécessaire |
Label Already Exists |
L'étiquette est déjà utilisée par une autre tâche | Modifiez l'étiquette et soumettez à nouveau |
Fail |
Échec du chargement | Consultez ErrorURL pour plus de détails, corrigez le problème et réessayez |
Vérifier les détails des erreurs
Lorsque Status est Fail, utilisez ErrorURL pour récupérer les détails concernant les lignes filtrées :
# View error rows directly
curl "<ErrorURL>"
# Save error rows to a local file
wget "<ErrorURL>"
Exemple :
wget "http://172.17.**.**:18040/api/_load_error_log?file=error_log_b74dccdcf0ceb4de_e82b2709c6c013ad"
Annuler une tâche de chargement
Il est impossible d'annuler manuellement les tâches Stream Load. Une tâche est automatiquement annulée en cas de délai d'expiration dépassé ou d'erreur d'importation.
Exemples
Charger des données CSV
Cet exemple charge data.csv dans example_table de la base de données load_test :
curl --location-trusted -u "root:" \
-H "Expect:100-continue" \
-H "label:label2" \
-H "column_separator:," \
-T data.csv -XPUT \
http://172.17.**.**:18030/api/load_test/example_table/_stream_load
Charger des données JSON
Cet exemple charge des enregistrements JSON depuis json.data :
curl --location-trusted -u "root:" \
-H "Expect:100-continue" \
-H "label:label2" \
-H "format:json" \
-T json.data -XPUT \
http://172.17.**.**:18030/api/load_test/example_table/_stream_load
Définir la tolérance aux erreurs
Pour autoriser jusqu'à 20 % des lignes à être filtrées en raison de problèmes de qualité des données :
curl --location-trusted -u "root:" \
-H "Expect:100-continue" \
-H "label:label3" \
-H "column_separator:," \
-H "max_filter_ratio:0.2" \
-T data.csv -XPUT \
http://172.17.**.**:18030/api/load_test/example_table/_stream_load
Filtrer les lignes par condition
Pour charger uniquement les lignes où k1 = 20180601 :
curl --location-trusted -u "root:" \
-H "Expect:100-continue" \
-H "label:label4" \
-H "column_separator:," \
-H "where:k1 = 20180601" \
-T data.csv -XPUT \
http://172.17.**.**:18030/api/load_test/example_table/_stream_load
Configurer le mappage des colonnes
Lorsque l'ordre des colonnes du fichier source ne correspond pas au schéma de la table, utilisez le paramètre columns pour spécifier le mappage.
Exemple 1 : La table possède les colonnes c1, c2, c3. Les colonnes du fichier source sont dans l'ordre c3, c2, c1 :
curl --location-trusted -u "root:" \
-H "Expect:100-continue" \
-H "label:label5" \
-H "column_separator:," \
-H "columns:c3, c2, c1" \
-T data.csv -XPUT \
http://172.17.**.**:18030/api/load_test/example_table/_stream_load
Exemple 2 : Le fichier source contient une colonne supplémentaire qui n'existe pas dans la table. Utilisez un espace réservé pour la colonne supplémentaire :
curl --location-trusted -u "root:" \
-H "Expect:100-continue" \
-H "label:label6" \
-H "column_separator:," \
-H "columns:c1, c2, c3, temp" \
-T data.csv -XPUT \
http://172.17.**.**:18030/api/load_test/example_table/_stream_load
Exemple 3 : La table possède les colonnes year, month, day. Le fichier source contient une seule colonne d'horodatage au format 2018-06-01 01:02:03. Utilisez des expressions de colonnes dérivées :
curl --location-trusted -u "root:" \
-H "Expect:100-continue" \
-H "label:label7" \
-H "column_separator:," \
-H "columns:col, year=year(col), month=month(col), day=day(col)" \
-T data.csv -XPUT \
http://172.17.**.**:18030/api/load_test/example_table/_stream_load
Charger des données compressées
Pour charger un fichier JSON compressé en LZ4 :
curl --location-trusted -u "root:" \
-H "Expect:100-continue" \
-H "label:label8" \
-H "format:json" \
-H "compression:lz4_frame" \
-T data.json.lz4 -XPUT \
http://172.17.**.**:18030/api/load_test/example_table/_stream_load
Exemple de bout en bout
Étape 1 : Créer la table cible
Connectez-vous au nœud maître du cluster StarRocks via SSH. Pour plus de détails, consultez Se connecter à un cluster.
-
Connectez-vous au cluster StarRocks à l'aide d'un client MySQL :
mysql -h127.0.0.1 -P 9030 -uroot -
Créez la base de données et la table :
CREATE DATABASE IF NOT EXISTS load_test; USE load_test; CREATE TABLE IF NOT EXISTS example_table ( id INT, name VARCHAR(50), age INT ) DUPLICATE KEY(id) DISTRIBUTED BY HASH(id) BUCKETS 3 PROPERTIES ( "replication_num" = "1" -- Set the number of replicas to 1. );Appuyez sur
Ctrl+Dpour quitter le client MySQL.
Étape 2 : Préparer les données sources
Préparer les données CSV
Données CSV — créez un fichier nommé data.csv :
id,name,age
1,Alice,25
2,Bob,30
3,Charlie,35
Préparer les données JSON
Données JSON — créez un fichier nommé json.data :
{"id":1,"name":"Emily","age":25}
{"id":2,"name":"Benjamin","age":35}
{"id":3,"name":"Olivia","age":28}
{"id":4,"name":"Alexander","age":60}
{"id":5,"name":"Ava","age":17}
Étape 3 : Exécuter la tâche de chargement
Importer des données CSV
Charger des données CSV :
curl --location-trusted -u "root:" \
-H "Expect:100-continue" \
-H "label:label1" \
-H "column_separator:," \
-T data.csv -XPUT \
http://172.17.**.**:18030/api/load_test/example_table/_stream_load
Importer des données JSON
Charger des données JSON :
curl --location-trusted -u "root:" \
-H "Expect:100-continue" \
-H "label:label2" \
-H "format:json" \
-T json.data -XPUT \
http://172.17.**.**:18030/api/load_test/example_table/_stream_load
Étapes suivantes
Pour des exemples d'intégration Java, consultez la démo stream_load.
Pour l'intégration Spark Streaming, consultez SparkStreaming vers StarRocks.