diff --git a/content/de/developer/integration/big-data/flink.md b/content/de/developer/integration/big-data/flink.md new file mode 100644 index 00000000..61f95c4b --- /dev/null +++ b/content/de/developer/integration/big-data/flink.md @@ -0,0 +1,278 @@ +--- +title: "Apache Flink" +description: "Lesen und schreiben Sie CSV-Daten im RustFS-Objektspeicher mit Apache Flink und dessen S3-Dateisystem-Plugin." +--- + +Diese Anleitung verbindet [Apache Flink](https://github.com/apache/flink) über Flinks S3-Dateisystem-Plugin (`flink-s3-fs-hadoop`) mit **RustFS**. Sie starten einen Session-Cluster mit Docker Compose, schreiben ein begrenztes Ergebnis als Batch in den Bucket und lesen es über Flink SQL zurück. Der Ablauf wurde mit `flink:1.20` und `rustfs/rustfs-x86-musl:v2.3.1` verifiziert. + +Sie benötigen Docker mit dem Compose-Plugin. Dieses Setup ist für lokale Integrationstests gedacht, nicht für den Produktivbetrieb. + +## Architektur + +```mermaid +flowchart LR + Job["Flink SQL job"] -->|"filesystem connector"| S3["S3 plugin (flink-s3-fs-hadoop)"] + S3 -->|"GET / PUT"| RustFS["RustFS :9000"] + Init["init-bucket job"] -->|"create my-bucket"| RustFS +``` + +Das Plugin `flink-s3-fs-hadoop` registriert das `s3://`-Schema für Flinks filesystem-Connector. Endpunkt, Path-Style-Adressierung, Plain HTTP und Anmeldeinformationen werden über `s3.*`-Properties in `flink-conf.yaml` konfiguriert (übergeben via `FLINK_PROPERTIES`). + +## 1. Projektdateien anlegen + +Erstellen Sie ein Arbeitsverzeichnis: + +```bash +mkdir rustfs-flink +cd rustfs-flink +``` + +Erstellen Sie eine Umgebungsdatei und ersetzen Sie beide Platzhalter für die Anmeldeinformationen: + +```ini title=".env" +RUSTFS_ACCESS_KEY= +RUSTFS_SECRET_KEY= +``` + +Verwenden Sie dedizierte Anmeldeinformationen für den Bucket. Committen Sie `.env` nicht in die Versionsverwaltung. + +Das S3-Plugin liegt im Image unter `/opt/flink/opt/` und muss geladen werden, indem es nach `/opt/flink/plugins/s3fs/` kopiert wird. Bereiten Sie ein lokales Verzeichnis dafür vor: + +```bash +mkdir -p s3fs +docker create --name flink-tmp flink:1.20 +docker cp flink-tmp:/opt/flink/opt/flink-s3-fs-hadoop-1.20.5.jar s3fs/ +docker rm flink-tmp +``` + +Erstellen Sie die Compose-Datei: + +```yaml title="compose.yaml" +services: + rustfs: + image: rustfs/rustfs-x86-musl:v2.3.1 + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + RUSTFS_VOLUMES: /data + RUSTFS_ADDRESS: ":9000" + RUSTFS_CONSOLE_ADDRESS: ":9001" + RUSTFS_CONSOLE_ENABLE: "true" + volumes: + - rustfs-data:/data + ports: + - "9000:9000" + - "9001:9001" + healthcheck: + test: ["CMD", "curl", "-sf", "http://127.0.0.1:9000/health"] + interval: 10s + timeout: 5s + retries: 6 + start_period: 10s + networks: + - flink + + create-bucket: + image: rustfs/rc:latest + depends_on: + rustfs: + condition: service_healthy + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + entrypoint: + - /bin/sh + - -c + - | + until /usr/bin/rc alias set rustfs http://rustfs:9000 "$${RUSTFS_ACCESS_KEY}" "$${RUSTFS_SECRET_KEY}"; do + echo "Waiting for RustFS..." + sleep 2 + done + /usr/bin/rc ls rustfs/my-bucket >/dev/null 2>&1 || /usr/bin/rc mb rustfs/my-bucket + networks: + - flink + + jobmanager: + image: flink:1.20 + command: jobmanager + environment: + FLINK_PROPERTIES: | + jobmanager.rpc.address: jobmanager + rest.address: jobmanager + rest.bind-address: 0.0.0.0 + s3.access-key: ${RUSTFS_ACCESS_KEY} + s3.secret-key: ${RUSTFS_SECRET_KEY} + s3.endpoint: http://rustfs:9000 + s3.path-style-access: true + volumes: + - ./s3fs:/opt/flink/plugins/s3fs:ro + ports: + - "8081:8081" + depends_on: + create-bucket: + condition: service_completed_successfully + networks: + - flink + + taskmanager: + image: flink:1.20 + command: taskmanager + environment: + FLINK_PROPERTIES: | + jobmanager.rpc.address: jobmanager + taskmanager.host: taskmanager + s3.access-key: ${RUSTFS_ACCESS_KEY} + s3.secret-key: ${RUSTFS_SECRET_KEY} + s3.endpoint: http://rustfs:9000 + s3.path-style-access: true + volumes: + - ./s3fs:/opt/flink/plugins/s3fs:ro + depends_on: + jobmanager: + condition: service_started + networks: + - flink + +networks: + flink: + +volumes: + rustfs-data: +``` + +Die Properties `s3.access-key`, `s3.secret-key`, `s3.endpoint` und `s3.path-style-access` konfigurieren das S3-Plugin auf JobManager und TaskManager. + +## 2. Bereitstellung starten + +Prüfen Sie die Compose-Datei, bevor Sie Container starten: + +```bash +docker compose config +``` + +Starten Sie die Dienste und warten Sie, bis die Bucket-Initialisierung abgeschlossen ist: + +```bash +docker compose up -d +docker compose ps -a +``` + +## 3. Ein Ergebnis nach RustFS schreiben + +Erstellen Sie den SQL-Job — Batch-Modus mit dem filesystem-Sink: + +```yaml title="batch.sql" +SET 'execution.runtime-mode' = 'batch'; + +CREATE TABLE sink ( + id INT, + payload STRING +) WITH ( + 'connector' = 'filesystem', + 'path' = 's3://my-bucket/flink-out/', + 'format' = 'csv' +); + +INSERT INTO sink + VALUES (1, 'alpha'), (2, 'bravo'), (3, 'charlie'), (4, 'delta'), (5, 'echo'); +``` + +Übermitteln Sie ihn über den SQL-Client im JobManager: + +```bash +docker compose exec jobmanager bash -c "/opt/flink/bin/sql-client.sh embedded -f /dev/stdin" < batch.sql +``` + +Der Job endet, sobald alle Zeilen geschrieben sind. + +## 4. Die Daten zurücklesen + +Erstellen Sie die Lesekonfiguration — der filesystem-Connector durchsucht das Präfix: + +```yaml title="read.sql" +CREATE TABLE readings ( + id INT, + payload STRING +) WITH ( + 'connector' = 'filesystem', + 'path' = 's3://my-bucket/flink-out/', + 'format' = 'csv' +); + +SET 'sql-client.execution.result-mode' = 'TABLEAU'; + +SELECT * FROM readings; +``` + +```bash +docker compose exec jobmanager bash -c "/opt/flink/bin/sql-client.sh embedded -f /dev/stdin" < read.sql +``` + +```text ++----+-------------+--------------------------------+ +| op | id | payload | ++----+-------------+--------------------------------+ +| +I | 1 | alpha | +| +I | 2 | bravo | +| +I | 3 | charlie | +| +I | 4 | delta | +| +I | 5 | echo | ++----+-------------+--------------------------------+ +``` + +## 5. Objekte in RustFS prüfen + +Listen Sie das Präfix über das Bucket-Initialisierungs-Image auf: + +```bash +docker compose run --rm --entrypoint /bin/sh create-bucket -c \ + '/usr/bin/rc alias set rustfs http://rustfs:9000 "$RUSTFS_ACCESS_KEY" "$RUSTFS_SECRET_KEY" >/dev/null && /usr/bin/rc ls rustfs/my-bucket/flink-out --recursive' +``` + +```text +[2026-09-21 01:34:52] 41 B flink-out/part-f759f9e8-3d1b-46a1-a92e-53e9b727e831-task-0-file-0 +``` + +Sie können das Präfix auch in der RustFS-Konsole anzeigen: + +![Die in der RustFS-Konsole gespeicherte Flink-Ausgabedatei](./images/rustfs-flink-objects.png) + +## 6. RustFS S3 Tables verwenden + +RustFS S3 Tables bietet einen integrierten Apache-Iceberg-REST-Katalog, sodass Flink einen Table-Bucket als verwaltetes Iceberg-Warehouse nutzen kann, während die Daten in RustFS bleiben. Aktivieren Sie einen Table-Bucket und richten Sie den REST-Katalog des Flink-Iceberg-Connectors auf RustFS, wie in [S3 Tables](/administration/data/s3-tables) beschrieben: Die REST-Katalog-URI ist `http://:9000/iceberg`, das Warehouse ist der Bucket-Name, und sowohl Kataloganfragen (AWS Signature Version 4, Signing-Name `s3`) als auch der S3-Dateizugriff verwenden Path-Style-Adressierung. + +Laut der S3-Tables-Support-Matrix validieren Sie die exakten Flink- und Iceberg-Versionen Ihrer Bereitstellung gegen den Katalog, bevor Sie diesen Pfad in Produktion übernehmen. + +## 7. Stack stoppen oder zurücksetzen + +Stoppen Sie die Container und behalten Sie das RustFS-Datenvolumen: + +```bash +docker compose down +``` + +Um die gespeicherten Dateien zu löschen und mit einem leeren RustFS-Volumen zu beginnen, fügen Sie ausdrücklich `--volumes` hinzu: + +```bash +docker compose down --volumes +``` + +## Fehlerbehebung + +### No AWS Credentials provided / AccessDenied bei Schreibvorgängen + +Das S3-Plugin liest seine Anmeldeinformationen aus den `s3.*`-Properties in `flink-conf.yaml`. Stellen Sie sicher, dass `s3.access-key`, `s3.secret-key`, `s3.endpoint` und `s3.path-style-access` in `FLINK_PROPERTIES` für **sowohl** den JobManager als auch den TaskManager gesetzt sind und dass das Plugin-jar auf beiden unter `/opt/flink/plugins/s3fs/` liegt. + +### Der TaskManager kann den Host `rustfs` nicht auflösen + +Alle Flink-Container und RustFS müssen sich ein Compose-Netzwerk teilen. Hängt RustFS an einem externen Netzwerk, verbinden Sie auch die Flink-Container damit, bevor Sie den Job übermitteln. + +### Ein Streaming-Schreibvorgang schlägt nach einem fehlgeschlagenen Versuch mit "Stream closed" fehl + +Die Wiederherstellung eines laufenden S3-Uploads nach einem Fehler kann den Writer in einen nicht wiederherstellbaren Zustand versetzen. Löschen Sie das Ausgabe-Präfix des Jobs im Bucket und übermitteln Sie den Job erneut, oder verwenden Sie für Einmal-Schreibvorgänge den Batch-Modus wie in dieser Anleitung. + +## Nächste Schritte + +- Lesen Sie die [S3-Kompatibilitätshinweise](/administration/protocols/s3), bevor Sie weitere S3-Operationen verwenden. +- Erstellen Sie dedizierte Produktions-Anmeldeinformationen mit dem [Access Key Management](/security-compliance/iam/access-token). +- Folgen Sie der [Apache-Flink-Dokumentation](https://nightlies.apache.org/flink/flink-docs-stable/) für filesystem-Connector-Optionen wie Partitionierung und Kompaktierung. diff --git a/content/de/developer/integration/big-data/images/rustfs-flink-objects.png b/content/de/developer/integration/big-data/images/rustfs-flink-objects.png new file mode 100644 index 00000000..caf4f204 Binary files /dev/null and b/content/de/developer/integration/big-data/images/rustfs-flink-objects.png differ diff --git a/content/de/developer/integration/big-data/images/rustfs-spark-objects.png b/content/de/developer/integration/big-data/images/rustfs-spark-objects.png new file mode 100644 index 00000000..f43539d7 Binary files /dev/null and b/content/de/developer/integration/big-data/images/rustfs-spark-objects.png differ diff --git a/content/de/developer/integration/big-data/images/rustfs-trino-objects.png b/content/de/developer/integration/big-data/images/rustfs-trino-objects.png new file mode 100644 index 00000000..3846a0fd Binary files /dev/null and b/content/de/developer/integration/big-data/images/rustfs-trino-objects.png differ diff --git a/content/de/developer/integration/big-data/index.md b/content/de/developer/integration/big-data/index.md index a36f346b..2ae0de3d 100644 --- a/content/de/developer/integration/big-data/index.md +++ b/content/de/developer/integration/big-data/index.md @@ -12,5 +12,8 @@ Use **RustFS** as the object storage layer for data analytics systems that suppo - [Milvus](./milvus.md) - [DuckDB](./duckdb.md) - [InfluxDB](./influxdb.md) +- [Spark](./spark.md) +- [Flink](./flink.md) +- [Trino](./trino.md) Keep application data in a dedicated bucket and prefix, and use credentials scoped to the required bucket operations. \ No newline at end of file diff --git a/content/de/developer/integration/big-data/meta.json b/content/de/developer/integration/big-data/meta.json index 69dca344..1d6f178f 100644 --- a/content/de/developer/integration/big-data/meta.json +++ b/content/de/developer/integration/big-data/meta.json @@ -1,10 +1,13 @@ { - "title": "Datenanalyse", + "title": "Data Analytics", "pages": [ "iceberg", "pyiceberg", "milvus", "duckdb", - "influxdb" + "influxdb", + "spark", + "flink", + "trino" ] } diff --git a/content/de/developer/integration/big-data/spark.md b/content/de/developer/integration/big-data/spark.md new file mode 100644 index 00000000..f03f33f3 --- /dev/null +++ b/content/de/developer/integration/big-data/spark.md @@ -0,0 +1,220 @@ +--- +title: "Apache Spark" +description: "Lesen und schreiben Sie Parquet-Daten im RustFS-Objektspeicher mit Apache Spark über den s3a-Connector." +--- + +Diese Anleitung verbindet [Apache Spark](https://github.com/apache/spark) über den `s3a`-Connector mit **RustFS**. Sie starten RustFS mit Docker Compose, führen einen Spark-Job aus, der ein Parquet-Dataset in den Bucket schreibt, lesen es zurück und prüfen die Objekte in RustFS. Der Ablauf wurde mit `apache/spark:3.5.6` (Hadoop 3.3.4 via `hadoop-aws`) und `rustfs/rustfs-x86-musl:v2.3.1` verifiziert. + +Sie benötigen Docker mit dem Compose-Plugin. Dieses Setup ist für lokale Integrationstests gedacht, nicht für den Produktivbetrieb. + +## Architektur + +```mermaid +flowchart LR + Spark["Spark driver + executors"] -->|"S3AFileSystem"| RustFS["RustFS :9000"] + Init["init-bucket job"] -->|"create my-bucket"| RustFS +``` + +Spark spricht über das `hadoop-aws`-S3A-Dateisystem mit RustFS. Die Connector-Einstellungen — Endpunkt, Path-Style-Adressierung, Plain HTTP und Anmeldeinformationen — werden als `spark.hadoop.fs.s3a.*`-Properties übergeben. + +## 1. Projektdateien anlegen + +Erstellen Sie ein Arbeitsverzeichnis: + +```bash +mkdir rustfs-spark +cd rustfs-spark +``` + +Erstellen Sie eine Umgebungsdatei und ersetzen Sie beide Platzhalter für die Anmeldeinformationen: + +```ini title=".env" +RUSTFS_ACCESS_KEY= +RUSTFS_SECRET_KEY= +``` + +Verwenden Sie dedizierte Anmeldeinformationen für den Bucket. Committen Sie `.env` nicht in die Versionsverwaltung. + +Erstellen Sie den Spark-Job: + +```python title="job.py" +from pyspark.sql import SparkSession + +spark = SparkSession.builder.appName("rustfs-spark-demo").getOrCreate() +spark.sparkContext.setLogLevel("WARN") + +spark.range(1000).withColumnRenamed("id", "num") \ + .write.mode("overwrite").parquet("s3a://my-bucket/spark-demo/events") + +back = spark.read.parquet("s3a://my-bucket/spark-demo/events") +print("ROWS_READ_BACK:", back.count()) +spark.stop() +``` + +Erstellen Sie die Compose-Datei: + +```yaml title="compose.yaml" +services: + rustfs: + image: rustfs/rustfs-x86-musl:v2.3.1 + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + RUSTFS_VOLUMES: /data + RUSTFS_ADDRESS: ":9000" + RUSTFS_CONSOLE_ADDRESS: ":9001" + RUSTFS_CONSOLE_ENABLE: "true" + volumes: + - rustfs-data:/data + ports: + - "9000:9000" + - "9001:9001" + healthcheck: + test: ["CMD", "curl", "-sf", "http://127.0.0.1:9000/health"] + interval: 10s + timeout: 5s + retries: 6 + start_period: 10s + networks: + - warehouse + + create-bucket: + image: rustfs/rc:latest + depends_on: + rustfs: + condition: service_healthy + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + entrypoint: + - /bin/sh + - -c + - | + until /usr/bin/rc alias set rustfs http://rustfs:9000 "$${RUSTFS_ACCESS_KEY}" "$${RUSTFS_SECRET_KEY}"; do + echo "Waiting for RustFS..." + sleep 2 + done + /usr/bin/rc ls rustfs/my-bucket >/dev/null 2>&1 || /usr/bin/rc mb rustfs/my-bucket + networks: + - warehouse + + spark: + image: apache/spark:3.5.6 + entrypoint: ["/opt/spark/bin/spark-submit"] + volumes: + - ./job.py:/job.py:ro + command: + - --conf + - spark.jars.ivy=/tmp/.ivy2 + - --packages + - org.apache.hadoop:hadoop-aws:3.3.4 + - --conf + - spark.hadoop.fs.s3a.endpoint=http://rustfs:9000 + - --conf + - spark.hadoop.fs.s3a.access.key=${RUSTFS_ACCESS_KEY} + - --conf + - spark.hadoop.fs.s3a.secret.key=${RUSTFS_SECRET_KEY} + - --conf + - spark.hadoop.fs.s3a.path.style.access=true + - --conf + - spark.hadoop.fs.s3a.connection.ssl.enabled=false + - --conf + - spark.hadoop.fs.s3a.impl=org.apache.hadoop.fs.s3a.S3AFileSystem + - /job.py + depends_on: + create-bucket: + condition: service_completed_successfully + networks: + - warehouse + +networks: + warehouse: + +volumes: + rustfs-data: +``` + +`--packages org.apache.hadoop:hadoop-aws:3.3.4` lädt den S3-Connector beim Start herunter; er muss zur Hadoop-Version des Spark-Images passen. `spark.jars.ivy=/tmp/.ivy2` verschiebt den Download-Cache in ein beschreibbares Verzeichnis. + +## 2. Speicher starten und Job ausführen + +Starten Sie die Speicherdienste: + +```bash +docker compose up -d +docker compose ps -a +``` + +Führen Sie den Spark-Job aus: + +```bash +docker compose run --rm spark +``` + +```text +ROWS_READ_BACK: 1000 +``` + +## 3. Objekte in RustFS prüfen + +Listen Sie das Dataset über das Bucket-Initialisierungs-Image auf: + +```bash +docker compose run --rm --entrypoint /bin/sh create-bucket -c \ + '/usr/bin/rc alias set rustfs http://rustfs:9000 "$RUSTFS_ACCESS_KEY" "$RUSTFS_SECRET_KEY" >/dev/null && /usr/bin/rc ls rustfs/my-bucket/spark-demo --recursive' +``` + +```text +[2026-09-21 00:58:22] 0 B spark-demo/events/_SUCCESS +[2026-09-21 00:58:21] 1.46 KiB spark-demo/events/part-00000-...-c000.snappy.parquet +[2026-09-21 00:58:21] 1.46 KiB spark-demo/events/part-00001-...-c000.snappy.parquet +[2026-09-21 00:58:21] 1.46 KiB spark-demo/events/part-00002-...-c000.snappy.parquet +``` + +Sie können das Präfix auch in der RustFS-Konsole anzeigen: + +![Spark-Parquet-Ausgabe in der RustFS-Konsole](./images/rustfs-spark-objects.png) + +## 4. RustFS S3 Tables verwenden + +RustFS S3 Tables bietet einen integrierten Apache-Iceberg-REST-Katalog, sodass Spark einen Table-Bucket als verwaltetes Iceberg-Warehouse nutzen kann, während die Daten in RustFS bleiben. Aktivieren Sie einen Table-Bucket und verbinden Sie den Iceberg-REST-Katalog von Spark wie in [S3 Tables](/administration/data/s3-tables) beschrieben: Die REST-Katalog-URI ist `http://:9000/iceberg`, das Warehouse ist der Bucket-Name, und sowohl Kataloganfragen (AWS Signature Version 4, Signing-Name `s3`) als auch der S3-Dateizugriff verwenden Path-Style-Adressierung. + +Laut der S3-Tables-Support-Matrix validieren Sie die exakten Spark- und Iceberg-Versionen Ihrer Bereitstellung gegen den Katalog, bevor Sie diesen Pfad in Produktion übernehmen. + +## 5. Stack stoppen oder zurücksetzen + +Stoppen Sie die Container und behalten Sie das RustFS-Datenvolumen: + +```bash +docker compose down +``` + +Um das Dataset zu löschen und mit einem leeren RustFS-Volumen zu beginnen, fügen Sie ausdrücklich `--volumes` hinzu: + +```bash +docker compose down --volumes +``` + +## Fehlerbehebung + +### NumberFormatException: For input string: "60s" + +Die `hadoop-aws`-Version passt nicht zur Hadoop-Version im Spark-Image. Spark-4.x-Images benötigen `hadoop-aws` 3.4.x; diese Anleitung pinnt `apache/spark:3.5.6` zusammen mit `hadoop-aws:3.3.4`. + +### Verbindungsfehler zu `rustfs:9000` + +`fs.s3a.endpoint` wird innerhalb des Compose-Netzwerks aufgelöst. Für einen Spark-Prozess auf dem Host verwenden Sie `http://localhost:9000`. + +### AccessDenied- oder 403-Antworten + +Stellen Sie sicher, dass die Connector-Einstellungen mit den RustFS-Anmeldeinformationen übereinstimmen und dass der Dienst `create-bucket` erfolgreich abgeschlossen wurde: + +```bash +docker compose logs create-bucket +``` + +## Nächste Schritte + +- Lesen Sie die [S3-Kompatibilitätshinweise](/administration/protocols/s3), bevor Sie weitere S3-Operationen verwenden. +- Erstellen Sie dedizierte Produktions-Anmeldeinformationen mit dem [Access Key Management](/security-compliance/iam/access-token). +- Folgen Sie der [Spark-Dokumentation](https://spark.apache.org/docs/latest/) für Structured Streaming und Datenquellenoptionen. diff --git a/content/de/developer/integration/big-data/trino.md b/content/de/developer/integration/big-data/trino.md new file mode 100644 index 00000000..b107416c --- /dev/null +++ b/content/de/developer/integration/big-data/trino.md @@ -0,0 +1,250 @@ +--- +title: "Trino" +description: "Abfragen von CSV- und Parquet-Daten im RustFS-Objektspeicher mit Trino und dem File-Metastore des hive-Connectors." +--- + +Diese Anleitung verbindet [Trino](https://github.com/trinodb/trino) — die verteilte SQL-Abfrage-Engine — über den hive-Connector mit seinem File-Metastore und das native S3-Dateisystem mit **RustFS**. Sie erstellen ein Schema und eine Tabelle, fügen Zeilen ein, lesen sie zurück und prüfen die Objekte in RustFS. Sowohl Tabellenmetadaten als auch Datendateien liegen in RustFS. Der Ablauf wurde mit `trinodb/trino:435` und `rustfs/rustfs-x86-musl:v2.3.1` verifiziert. + +Sie benötigen Docker mit dem Compose-Plugin. Dieses Setup ist für lokale Integrationstests gedacht, nicht für den Produktivbetrieb. + +## Architektur + +```mermaid +flowchart LR + Client["trino CLI"] -->|"SQL"| Trino["Trino :8080"] + Trino -->|"metadata JSON"| RustFS["RustFS :9000"] + Trino -->|"data files"| RustFS + Init["init-bucket job"] -->|"create my-bucket"| RustFS +``` + +Der hive-Connector mit `hive.metastore=file` hält Schema- und Tabellenmetadaten als JSON-Objekte unter dem Katalogverzeichnis, und das native S3-Dateisystem (`fs.s3.enabled`) speichert Metadaten und Datendateien in RustFS mit Path-Style-Adressierung über Plain HTTP. + +## 1. Projektdateien anlegen + +Erstellen Sie ein Arbeitsverzeichnis: + +```bash +mkdir rustfs-trino +cd rustfs-trino +``` + +Erstellen Sie eine Umgebungsdatei und ersetzen Sie beide Platzhalter für die Anmeldeinformationen: + +```ini title=".env" +RUSTFS_ACCESS_KEY= +RUSTFS_SECRET_KEY= +``` + +Verwenden Sie dedizierte Anmeldeinformationen für den Bucket. Committen Sie `.env` nicht in die Versionsverwaltung. + +Erstellen Sie die Katalog-Konfiguration für Trino: + +```ini title="hive.properties" +connector.name=hive +hive.metastore=file +hive.metastore.catalog.dir=s3://my-bucket/trino-metastore +fs.s3.enabled=true +s3.endpoint=http://rustfs:9000 +s3.region=us-east-1 +s3.path-style-access=true +s3.aws-access-key= +s3.aws-secret-key= +``` + +`hive.metastore.catalog.dir` richtet den File-Metastore in den Bucket, sodass Metadaten und Daten beide in RustFS liegen. `fs.s3.enabled` aktiviert das native S3-Dateisystem; `s3.path-style-access` ist für den Container-Netzwerk-Endpunkt erforderlich. + +Erstellen Sie die Compose-Datei: + +```yaml title="compose.yaml" +services: + rustfs: + image: rustfs/rustfs-x86-musl:v2.3.1 + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + RUSTFS_VOLUMES: /data + RUSTFS_ADDRESS: ":9000" + RUSTFS_CONSOLE_ADDRESS: ":9001" + RUSTFS_CONSOLE_ENABLE: "true" + volumes: + - rustfs-data:/data + ports: + - "9000:9000" + - "9001:9001" + healthcheck: + test: ["CMD", "curl", "-sf", "http://127.0.0.1:9000/health"] + interval: 10s + timeout: 5s + retries: 6 + start_period: 10s + networks: + - warehouse + + create-bucket: + image: rustfs/rc:latest + depends_on: + rustfs: + condition: service_healthy + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + entrypoint: + - /bin/sh + - -c + - | + until /usr/bin/rc alias set rustfs http://rustfs:9000 "$${RUSTFS_ACCESS_KEY}" "$${RUSTFS_SECRET_KEY}"; do + echo "Waiting for RustFS..." + sleep 2 + done + /usr/bin/rc ls rustfs/my-bucket >/dev/null 2>&1 || /usr/bin/rc mb rustfs/my-bucket + networks: + - warehouse + + trino: + image: trinodb/trino:435 + volumes: + - ./hive.properties:/etc/trino/catalog/hive.properties:ro + - metastore-data:/data/metastore + depends_on: + create-bucket: + condition: service_completed_successfully + networks: + - warehouse + +networks: + warehouse: + +volumes: + rustfs-data: + metastore-data: +``` + +## 2. Bereitstellung starten + +Prüfen Sie die Compose-Datei, bevor Sie Container starten: + +```bash +docker compose config +``` + +Starten Sie die Dienste und warten Sie, bis die Bucket-Initialisierung abgeschlossen ist: + +```bash +docker compose up -d +docker compose ps -a +``` + +Trino läuft, wenn das Server-Log `SERVER STARTED` meldet. Der Container läuft als Benutzer `trino` (uid 1000); stellen Sie sicher, dass das Metastore-Volume beschreibbar ist: + +```bash +docker compose exec trino id +docker compose exec trino ls -la /data/metastore +``` + +## 3. Schema und Tabelle erstellen + +Erstellen Sie das Schema ohne explizite Location — Trino legt es unter dem Katalogverzeichnis in RustFS ab: + +```bash +docker compose exec trino trino --execute \ + "CREATE SCHEMA hive.demo" +``` + +Erstellen Sie eine Tabelle und fügen Sie fünf Zeilen ein: + +```bash +docker compose exec trino trino --execute \ + "CREATE TABLE hive.demo.events (id bigint, label varchar) WITH (format = 'parquet')" + +docker compose exec trino trino --execute \ + "INSERT INTO hive.demo.events VALUES (1,'alpha'),(2,'bravo'),(3,'charlie'),(4,'delta'),(5,'echo')" +``` + +```text +INSERT: 5 rows +``` + +## 4. Die Daten abfragen + +Lesen Sie die Zeilen zurück: + +```bash +docker compose exec trino trino --execute \ + "SELECT * FROM hive.demo.events ORDER BY id" +``` + +```text +"1","alpha" +"2","bravo" +"3","charlie" +"4","delta" +"5","echo" +``` + +## 5. Objekte in RustFS prüfen + +Listen Sie das Metastore-Präfix über das Bucket-Initialisierungs-Image auf: + +```bash +docker compose run --rm --entrypoint /bin/sh create-bucket -c \ + '/usr/bin/rc alias set rustfs http://rustfs:9000 "$RUSTFS_ACCESS_KEY" "$RUSTFS_SECRET_KEY" >/dev/null && /usr/bin/rc ls rustfs/my-bucket/trino-metastore --recursive' +``` + +```text +[2026-09-21 01:55:16] 155 B trino-metastore/.demo.trinoSchema +[2026-09-21 01:55:19] 474 B trino-metastore/demo/events/.trinoPermissions/user_trino +[2026-09-21 01:55:25] 1007 B trino-metastore/demo/events/.trinoSchema +[2026-09-21 01:55:25] 432 B trino-metastore/demo/events/20260921_..._cb761cec-...parquet +``` + +Sie können das Präfix auch in der RustFS-Konsole anzeigen: + +![Die Trino-Metadaten- und Datenobjekte in der RustFS-Konsole](./images/rustfs-trino-objects.png) + +## 6. RustFS S3 Tables verwenden + +RustFS S3 Tables bietet einen integrierten Apache-Iceberg-REST-Katalog, sodass Trino einen Table-Bucket als verwaltetes Iceberg-Warehouse nutzen kann, während die Daten in RustFS bleiben. Aktivieren Sie einen Table-Bucket und verbinden Sie den Iceberg-Connector von Trino mit dem REST-Katalog, wie in [S3 Tables](/administration/data/s3-tables) beschrieben: Die REST-Katalog-URI ist `http://:9000/iceberg`, das Warehouse ist der Bucket-Name, und sowohl Kataloganfragen (AWS Signature Version 4, Signing-Name `s3`) als auch der S3-Dateizugriff verwenden Path-Style-Adressierung. + +Laut der S3-Tables-Support-Matrix wurde Trino nur als Read-only-Probe gegen den Katalog geprüft; validieren Sie die Schreibkompatibilität und die exakte Trino-Version Ihrer Bereitstellung, bevor Sie diesen Pfad in Produktion übernehmen. + +## 7. Stack stoppen oder zurücksetzen + +Stoppen Sie die Container und behalten Sie das RustFS-Datenvolumen: + +```bash +docker compose down +``` + +Um die gespeicherten Metadaten und Daten zu löschen und mit einem leeren RustFS-Volumen zu beginnen, fügen Sie ausdrücklich `--volumes` hinzu: + +```bash +docker compose down --volumes +``` + +## Fehlerbehebung + +### Konfigurationsfehler für `fs.native-s3.enabled` oder `fs.s3.enabled` + +Der Property-Name des nativen S3-Dateisystems hat sich zwischen Trino-Versionen geändert: Trino 435 verwendet `fs.native-s3.enabled`, neuere Releases `fs.s3.enabled`. Diese Anleitung pinnt `trinodb/trino:435`, daher verwenden Sie `fs.native-s3.enabled`. + +### "Table directory must be ..." beim Erstellen einer Tabelle + +Mit dem File-Metastore müssen Tabellen-Locations unter `hive.metastore.catalog.dir` bleiben. Erstellen Sie das Schema ohne explizite Location, oder zeigen Sie die Schema-Location auf ein Verzeichnis innerhalb desselben Bucket-Präfixes. + +### Hive CSV storage format only supports VARCHAR + +Das CSV-Format lehnt Nicht-String-Spalten ab. Verwenden Sie für typisierte Tabellen `format = 'parquet'` wie in dieser Anleitung. + +### AccessDenied- oder 403-Antworten + +Stellen Sie sicher, dass die Anmeldeinformationen in `hive.properties` mit den RustFS-Anmeldeinformationen übereinstimmen und dass der Dienst `create-bucket` erfolgreich abgeschlossen wurde: + +```bash +docker compose logs create-bucket +``` + +## Nächste Schritte + +- Lesen Sie die [S3-Kompatibilitätshinweise](/administration/protocols/s3), bevor Sie weitere S3-Operationen verwenden. +- Erstellen Sie dedizierte Produktions-Anmeldeinformationen mit dem [Access Key Management](/security-compliance/iam/access-token). +- Folgen Sie der [Trino-Dokumentation](https://trino.io/docs/current/), um BI-Tools anzubinden und weitere Objektspeicher-Kataloge hinzuzufügen. diff --git a/content/en/developer/integration/big-data/flink.md b/content/en/developer/integration/big-data/flink.md new file mode 100644 index 00000000..586b1f68 --- /dev/null +++ b/content/en/developer/integration/big-data/flink.md @@ -0,0 +1,278 @@ +--- +title: "Apache Flink" +description: "Read from and write CSV data stored in RustFS object storage with Apache Flink and its S3 filesystem plugin." +--- + +This guide connects [Apache Flink](https://github.com/apache/flink) to **RustFS** through Flink's S3 filesystem plugin (`flink-s3-fs-hadoop`). You will start a session cluster with Docker Compose, write a bounded result set to the bucket in batch mode, and read it back through Flink SQL. The workflow was verified with `flink:1.20` and `rustfs/rustfs-x86-musl:v2.3.1`. + +You need Docker with the Compose plugin. This deployment is intended for local integration testing, not production. + +## Architecture + +```mermaid +flowchart LR + Job["Flink SQL job"] -->|"filesystem connector"| S3["S3 plugin (flink-s3-fs-hadoop)"] + S3 -->|"GET / PUT"| RustFS["RustFS :9000"] + Init["init-bucket job"] -->|"create my-bucket"| RustFS +``` + +The `flink-s3-fs-hadoop` plugin registers the `s3://` scheme for Flink's filesystem connector. Endpoint, path-style addressing, plain HTTP, and credentials are configured through `s3.*` properties in `flink-conf.yaml` (passed via `FLINK_PROPERTIES`). + +## 1. Create the project files + +Create a working directory: + +```bash +mkdir rustfs-flink +cd rustfs-flink +``` + +Create an environment file and replace both credential placeholders: + +```ini title=".env" +RUSTFS_ACCESS_KEY= +RUSTFS_SECRET_KEY= +``` + +Use dedicated credentials for the bucket. Do not commit `.env` to source control. + +The S3 plugin ships inside the image under `/opt/flink/opt/` and must be copied to `/opt/flink/plugins/s3fs/` to load. Prepare a local directory for it: + +```bash +mkdir -p s3fs +docker create --name flink-tmp flink:1.20 +docker cp flink-tmp:/opt/flink/opt/flink-s3-fs-hadoop-1.20.5.jar s3fs/ +docker rm flink-tmp +``` + +Create the Compose file: + +```yaml title="compose.yaml" +services: + rustfs: + image: rustfs/rustfs-x86-musl:v2.3.1 + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + RUSTFS_VOLUMES: /data + RUSTFS_ADDRESS: ":9000" + RUSTFS_CONSOLE_ADDRESS: ":9001" + RUSTFS_CONSOLE_ENABLE: "true" + volumes: + - rustfs-data:/data + ports: + - "9000:9000" + - "9001:9001" + healthcheck: + test: ["CMD", "curl", "-sf", "http://127.0.0.1:9000/health"] + interval: 10s + timeout: 5s + retries: 6 + start_period: 10s + networks: + - flink + + create-bucket: + image: rustfs/rc:latest + depends_on: + rustfs: + condition: service_healthy + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + entrypoint: + - /bin/sh + - -c + - | + until /usr/bin/rc alias set rustfs http://rustfs:9000 "$${RUSTFS_ACCESS_KEY}" "$${RUSTFS_SECRET_KEY}"; do + echo "Waiting for RustFS..." + sleep 2 + done + /usr/bin/rc ls rustfs/my-bucket >/dev/null 2>&1 || /usr/bin/rc mb rustfs/my-bucket + networks: + - flink + + jobmanager: + image: flink:1.20 + command: jobmanager + environment: + FLINK_PROPERTIES: | + jobmanager.rpc.address: jobmanager + rest.address: jobmanager + rest.bind-address: 0.0.0.0 + s3.access-key: ${RUSTFS_ACCESS_KEY} + s3.secret-key: ${RUSTFS_SECRET_KEY} + s3.endpoint: http://rustfs:9000 + s3.path-style-access: true + volumes: + - ./s3fs:/opt/flink/plugins/s3fs:ro + ports: + - "8081:8081" + depends_on: + create-bucket: + condition: service_completed_successfully + networks: + - flink + + taskmanager: + image: flink:1.20 + command: taskmanager + environment: + FLINK_PROPERTIES: | + jobmanager.rpc.address: jobmanager + taskmanager.host: taskmanager + s3.access-key: ${RUSTFS_ACCESS_KEY} + s3.secret-key: ${RUSTFS_SECRET_KEY} + s3.endpoint: http://rustfs:9000 + s3.path-style-access: true + volumes: + - ./s3fs:/opt/flink/plugins/s3fs:ro + depends_on: + jobmanager: + condition: service_started + networks: + - flink + +networks: + flink: + +volumes: + rustfs-data: +``` + +The `s3.access-key`, `s3.secret-key`, `s3.endpoint`, and `s3.path-style-access` properties configure the S3 plugin on both the JobManager and the TaskManager. + +## 2. Start the deployment + +Resolve the Compose file before starting containers: + +```bash +docker compose config +``` + +Start the services and wait for the bucket initializer to finish: + +```bash +docker compose up -d +docker compose ps -a +``` + +## 3. Write a result set to RustFS + +Create the SQL job — batch mode with the filesystem sink: + +```yaml title="batch.sql" +SET 'execution.runtime-mode' = 'batch'; + +CREATE TABLE sink ( + id INT, + payload STRING +) WITH ( + 'connector' = 'filesystem', + 'path' = 's3://my-bucket/flink-out/', + 'format' = 'csv' +); + +INSERT INTO sink + VALUES (1, 'alpha'), (2, 'bravo'), (3, 'charlie'), (4, 'delta'), (5, 'echo'); +``` + +Submit it through the SQL client inside the JobManager: + +```bash +docker compose exec jobmanager bash -c "/opt/flink/bin/sql-client.sh embedded -f /dev/stdin" < batch.sql +``` + +The job finishes when all rows are written. + +## 4. Read the data back + +Create the read query — the filesystem connector scans the prefix: + +```yaml title="read.sql" +CREATE TABLE readings ( + id INT, + payload STRING +) WITH ( + 'connector' = 'filesystem', + 'path' = 's3://my-bucket/flink-out/', + 'format' = 'csv' +); + +SET 'sql-client.execution.result-mode' = 'TABLEAU'; + +SELECT * FROM readings; +``` + +```bash +docker compose exec jobmanager bash -c "/opt/flink/bin/sql-client.sh embedded -f /dev/stdin" < read.sql +``` + +```text ++----+-------------+--------------------------------+ +| op | id | payload | ++----+-------------+--------------------------------+ +| +I | 1 | alpha | +| +I | 2 | bravo | +| +I | 3 | charlie | +| +I | 4 | delta | +| +I | 5 | echo | ++----+-------------+--------------------------------+ +``` + +## 5. Verify objects in RustFS + +List the prefix through the bucket-initializer image: + +```bash +docker compose run --rm --entrypoint /bin/sh create-bucket -c \ + '/usr/bin/rc alias set rustfs http://rustfs:9000 "$RUSTFS_ACCESS_KEY" "$RUSTFS_SECRET_KEY" >/dev/null && /usr/bin/rc ls rustfs/my-bucket/flink-out --recursive' +``` + +```text +[2026-09-21 01:34:52] 41 B flink-out/part-f759f9e8-3d1b-46a1-a92e-53e9b727e831-task-0-file-0 +``` + +You can also browse the prefix in the RustFS Console: + +![The Flink output file stored in the RustFS Console](./images/rustfs-flink-objects.png) + +## 6. Use RustFS S3 Tables + +RustFS S3 Tables provides a built-in Apache Iceberg REST catalog, so Flink can treat a table bucket as a managed Iceberg warehouse while the data stays in RustFS. Enable a table bucket and point the Flink Iceberg connector's REST catalog at RustFS as described in [S3 Tables](/administration/data/s3-tables): the REST catalog URI is `http://:9000/iceberg`, the warehouse is the bucket name, and both catalog requests (AWS Signature Version 4, signing name `s3`) and S3 file access use path-style addressing. + +Per the S3 Tables support matrix, validate the exact Flink and Iceberg versions you deploy against the catalog before adopting this path in production. + +## 7. Stop or reset the stack + +Stop the containers while keeping the RustFS data volume: + +```bash +docker compose down +``` + +To delete the stored files and start from an empty RustFS volume, explicitly include `--volumes`: + +```bash +docker compose down --volumes +``` + +## Troubleshooting + +### No AWS Credentials provided / AccessDenied on writes + +The S3 plugin reads its credentials from the `s3.*` properties in `flink-conf.yaml`. Confirm that `s3.access-key`, `s3.secret-key`, `s3.endpoint`, and `s3.path-style-access` are present in `FLINK_PROPERTIES` for **both** the JobManager and the TaskManager, and that the plugin jar exists in `/opt/flink/plugins/s3fs/` on each. + +### Cannot resolve the `rustfs` host from the TaskManager + +All Flink containers and RustFS must share a Compose network. If you attach RustFS to an external network, connect the Flink containers to it as well before submitting the job. + +### A streaming write fails with "Stream closed" after a failed attempt + +Recovering an in-progress S3 upload after a failure can leave the writer in an unrecoverable state. Delete the job's output prefix in the bucket and resubmit the job, or use batch mode as shown in this guide for one-shot writes. + +## Next steps + +- Review [S3 compatibility notes](/administration/protocols/s3) before adopting additional S3 operations. +- Create dedicated production credentials with [Access Key Management](/security-compliance/iam/access-token). +- Follow the [Apache Flink documentation](https://nightlies.apache.org/flink/flink-docs-stable/) for filesystem connector options such as partitioning and compaction. diff --git a/content/en/developer/integration/big-data/images/rustfs-flink-objects.png b/content/en/developer/integration/big-data/images/rustfs-flink-objects.png new file mode 100644 index 00000000..caf4f204 Binary files /dev/null and b/content/en/developer/integration/big-data/images/rustfs-flink-objects.png differ diff --git a/content/en/developer/integration/big-data/images/rustfs-spark-objects.png b/content/en/developer/integration/big-data/images/rustfs-spark-objects.png new file mode 100644 index 00000000..f43539d7 Binary files /dev/null and b/content/en/developer/integration/big-data/images/rustfs-spark-objects.png differ diff --git a/content/en/developer/integration/big-data/images/rustfs-trino-objects.png b/content/en/developer/integration/big-data/images/rustfs-trino-objects.png new file mode 100644 index 00000000..3846a0fd Binary files /dev/null and b/content/en/developer/integration/big-data/images/rustfs-trino-objects.png differ diff --git a/content/en/developer/integration/big-data/index.md b/content/en/developer/integration/big-data/index.md index 055a3a41..2325206d 100644 --- a/content/en/developer/integration/big-data/index.md +++ b/content/en/developer/integration/big-data/index.md @@ -12,5 +12,8 @@ Use **RustFS** as the object storage layer for data analytics systems that suppo - [Milvus](./milvus.md) - [DuckDB](./duckdb.md) - [InfluxDB](./influxdb.md) +- [Spark](./spark.md) +- [Flink](./flink.md) +- [Trino](./trino.md) Keep application data in a dedicated bucket and prefix, and use credentials scoped to the required bucket operations. \ No newline at end of file diff --git a/content/en/developer/integration/big-data/meta.json b/content/en/developer/integration/big-data/meta.json index f4b58a3e..1d6f178f 100644 --- a/content/en/developer/integration/big-data/meta.json +++ b/content/en/developer/integration/big-data/meta.json @@ -5,6 +5,9 @@ "pyiceberg", "milvus", "duckdb", - "influxdb" + "influxdb", + "spark", + "flink", + "trino" ] } diff --git a/content/en/developer/integration/big-data/spark.md b/content/en/developer/integration/big-data/spark.md new file mode 100644 index 00000000..bd435955 --- /dev/null +++ b/content/en/developer/integration/big-data/spark.md @@ -0,0 +1,220 @@ +--- +title: "Apache Spark" +description: "Read and write Parquet data stored in RustFS object storage with Apache Spark over the s3a connector." +--- + +This guide connects [Apache Spark](https://github.com/apache/spark) to **RustFS** through the `s3a` connector. You will start RustFS with Docker Compose, run a Spark job that writes a Parquet dataset to the bucket, read it back, and verify the objects in RustFS. The workflow was verified with `apache/spark:3.5.6` (Hadoop 3.3.4 via `hadoop-aws`) and `rustfs/rustfs-x86-musl:v2.3.1`. + +You need Docker with the Compose plugin. This deployment is intended for local integration testing, not production. + +## Architecture + +```mermaid +flowchart LR + Spark["Spark driver + executors"] -->|"S3AFileSystem"| RustFS["RustFS :9000"] + Init["init-bucket job"] -->|"create my-bucket"| RustFS +``` + +Spark talks to RustFS through the `hadoop-aws` S3A filesystem. The connector settings — endpoint, path-style addressing, plain HTTP, and credentials — are passed as `spark.hadoop.fs.s3a.*` properties. + +## 1. Create the project files + +Create a working directory: + +```bash +mkdir rustfs-spark +cd rustfs-spark +``` + +Create an environment file and replace both credential placeholders: + +```ini title=".env" +RUSTFS_ACCESS_KEY= +RUSTFS_SECRET_KEY= +``` + +Use dedicated credentials for the bucket. Do not commit `.env` to source control. + +Create the Spark job: + +```python title="job.py" +from pyspark.sql import SparkSession + +spark = SparkSession.builder.appName("rustfs-spark-demo").getOrCreate() +spark.sparkContext.setLogLevel("WARN") + +spark.range(1000).withColumnRenamed("id", "num") \ + .write.mode("overwrite").parquet("s3a://my-bucket/spark-demo/events") + +back = spark.read.parquet("s3a://my-bucket/spark-demo/events") +print("ROWS_READ_BACK:", back.count()) +spark.stop() +``` + +Create the Compose file: + +```yaml title="compose.yaml" +services: + rustfs: + image: rustfs/rustfs-x86-musl:v2.3.1 + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + RUSTFS_VOLUMES: /data + RUSTFS_ADDRESS: ":9000" + RUSTFS_CONSOLE_ADDRESS: ":9001" + RUSTFS_CONSOLE_ENABLE: "true" + volumes: + - rustfs-data:/data + ports: + - "9000:9000" + - "9001:9001" + healthcheck: + test: ["CMD", "curl", "-sf", "http://127.0.0.1:9000/health"] + interval: 10s + timeout: 5s + retries: 6 + start_period: 10s + networks: + - warehouse + + create-bucket: + image: rustfs/rc:latest + depends_on: + rustfs: + condition: service_healthy + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + entrypoint: + - /bin/sh + - -c + - | + until /usr/bin/rc alias set rustfs http://rustfs:9000 "$${RUSTFS_ACCESS_KEY}" "$${RUSTFS_SECRET_KEY}"; do + echo "Waiting for RustFS..." + sleep 2 + done + /usr/bin/rc ls rustfs/my-bucket >/dev/null 2>&1 || /usr/bin/rc mb rustfs/my-bucket + networks: + - warehouse + + spark: + image: apache/spark:3.5.6 + entrypoint: ["/opt/spark/bin/spark-submit"] + volumes: + - ./job.py:/job.py:ro + command: + - --conf + - spark.jars.ivy=/tmp/.ivy2 + - --packages + - org.apache.hadoop:hadoop-aws:3.3.4 + - --conf + - spark.hadoop.fs.s3a.endpoint=http://rustfs:9000 + - --conf + - spark.hadoop.fs.s3a.access.key=${RUSTFS_ACCESS_KEY} + - --conf + - spark.hadoop.fs.s3a.secret.key=${RUSTFS_SECRET_KEY} + - --conf + - spark.hadoop.fs.s3a.path.style.access=true + - --conf + - spark.hadoop.fs.s3a.connection.ssl.enabled=false + - --conf + - spark.hadoop.fs.s3a.impl=org.apache.hadoop.fs.s3a.S3AFileSystem + - /job.py + depends_on: + create-bucket: + condition: service_completed_successfully + networks: + - warehouse + +networks: + warehouse: + +volumes: + rustfs-data: +``` + +`--packages org.apache.hadoop:hadoop-aws:3.3.4` downloads the S3 connector at launch; it must match the Hadoop version bundled with the Spark image. `spark.jars.ivy=/tmp/.ivy2` moves the download cache to a writable directory. + +## 2. Start RustFS and run the job + +Start the storage services: + +```bash +docker compose up -d +docker compose ps -a +``` + +Run the Spark job: + +```bash +docker compose run --rm spark +``` + +```text +ROWS_READ_BACK: 1000 +``` + +## 3. Verify objects in RustFS + +List the dataset through the bucket-initializer image: + +```bash +docker compose run --rm --entrypoint /bin/sh create-bucket -c \ + '/usr/bin/rc alias set rustfs http://rustfs:9000 "$RUSTFS_ACCESS_KEY" "$RUSTFS_SECRET_KEY" >/dev/null && /usr/bin/rc ls rustfs/my-bucket/spark-demo --recursive' +``` + +```text +[2026-09-21 00:58:22] 0 B spark-demo/events/_SUCCESS +[2026-09-21 00:58:21] 1.46 KiB spark-demo/events/part-00000-...-c000.snappy.parquet +[2026-09-21 00:58:21] 1.46 KiB spark-demo/events/part-00001-...-c000.snappy.parquet +[2026-09-21 00:58:21] 1.46 KiB spark-demo/events/part-00002-...-c000.snappy.parquet +``` + +You can also browse the prefix in the RustFS Console: + +![Spark Parquet output stored in the RustFS Console](./images/rustfs-spark-objects.png) + +## 4. Use RustFS S3 Tables + +RustFS S3 Tables provides a built-in Apache Iceberg REST catalog, so Spark can treat a table bucket as a managed Iceberg warehouse while the data stays in RustFS. Enable a table bucket and connect Spark's Iceberg REST catalog as described in [S3 Tables](/administration/data/s3-tables): the REST catalog URI is `http://:9000/iceberg`, the warehouse is the bucket name, and both catalog requests (AWS Signature Version 4, signing name `s3`) and S3 file access use path-style addressing. + +Per the S3 Tables support matrix, validate the exact Spark and Iceberg versions you deploy against the catalog before adopting this path in production. + +## 5. Stop or reset the stack + +Stop the containers while keeping the RustFS data volume: + +```bash +docker compose down +``` + +To delete the dataset and start from an empty RustFS volume, explicitly include `--volumes`: + +```bash +docker compose down --volumes +``` + +## Troubleshooting + +### NumberFormatException: For input string: "60s" + +The `hadoop-aws` version does not match the Hadoop version inside the Spark image. Spark 4.x images need `hadoop-aws` 3.4.x; this guide pins `apache/spark:3.5.6` together with `hadoop-aws:3.3.4`. + +### Connection failures to `rustfs:9000` + +`fs.s3a.endpoint` is resolved inside the Compose network. From a Spark process running on the host, use `http://localhost:9000` instead. + +### AccessDenied or 403 responses + +Confirm that the connector settings match the RustFS credentials and that the `create-bucket` service completed successfully: + +```bash +docker compose logs create-bucket +``` + +## Next steps + +- Review [S3 compatibility notes](/administration/protocols/s3) before adopting additional S3 operations. +- Create dedicated production credentials with [Access Key Management](/security-compliance/iam/access-token). +- Follow the [Spark documentation](https://spark.apache.org/docs/latest/) for structured streaming and data source options. diff --git a/content/en/developer/integration/big-data/trino.md b/content/en/developer/integration/big-data/trino.md new file mode 100644 index 00000000..67c08fce --- /dev/null +++ b/content/en/developer/integration/big-data/trino.md @@ -0,0 +1,250 @@ +--- +title: "Trino" +description: "Query CSV and Parquet data stored in RustFS object storage with Trino and the hive connector's file metastore." +--- + +This guide connects [Trino](https://github.com/trinodb/trino) — the distributed SQL query engine — to **RustFS** through the hive connector with its file-based metastore and the native S3 filesystem. You will create a schema and a table, insert rows, read them back, and verify the objects in RustFS. Both the table metadata and the data files live in RustFS. The workflow was verified with `trinodb/trino:435` and `rustfs/rustfs-x86-musl:v2.3.1`. + +You need Docker with the Compose plugin. This deployment is intended for local integration testing, not production. + +## Architecture + +```mermaid +flowchart LR + Client["trino CLI"] -->|"SQL"| Trino["Trino :8080"] + Trino -->|"metadata JSON"| RustFS["RustFS :9000"] + Trino -->|"data files"| RustFS + Init["init-bucket job"] -->|"create my-bucket"| RustFS +``` + +The hive connector with `hive.metastore=file` keeps schema and table metadata as JSON objects under the catalog directory, and the native S3 filesystem (`fs.s3.enabled`) stores both metadata and data files in RustFS with path-style addressing over plain HTTP. + +## 1. Create the project files + +Create a working directory: + +```bash +mkdir rustfs-trino +cd rustfs-trino +``` + +Create an environment file and replace both credential placeholders: + +```ini title=".env" +RUSTFS_ACCESS_KEY= +RUSTFS_SECRET_KEY= +``` + +Use dedicated credentials for the bucket. Do not commit `.env` to source control. + +Create the catalog configuration for Trino: + +```ini title="hive.properties" +connector.name=hive +hive.metastore=file +hive.metastore.catalog.dir=s3://my-bucket/trino-metastore +fs.s3.enabled=true +s3.endpoint=http://rustfs:9000 +s3.region=us-east-1 +s3.path-style-access=true +s3.aws-access-key= +s3.aws-secret-key= +``` + +`hive.metastore.catalog.dir` points the file metastore into the bucket, so metadata and data both live in RustFS. `fs.s3.enabled` activates the native S3 filesystem; `s3.path-style-access` is required for the container-network endpoint. + +Create the Compose file: + +```yaml title="compose.yaml" +services: + rustfs: + image: rustfs/rustfs-x86-musl:v2.3.1 + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + RUSTFS_VOLUMES: /data + RUSTFS_ADDRESS: ":9000" + RUSTFS_CONSOLE_ADDRESS: ":9001" + RUSTFS_CONSOLE_ENABLE: "true" + volumes: + - rustfs-data:/data + ports: + - "9000:9000" + - "9001:9001" + healthcheck: + test: ["CMD", "curl", "-sf", "http://127.0.0.1:9000/health"] + interval: 10s + timeout: 5s + retries: 6 + start_period: 10s + networks: + - warehouse + + create-bucket: + image: rustfs/rc:latest + depends_on: + rustfs: + condition: service_healthy + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + entrypoint: + - /bin/sh + - -c + - | + until /usr/bin/rc alias set rustfs http://rustfs:9000 "$${RUSTFS_ACCESS_KEY}" "$${RUSTFS_SECRET_KEY}"; do + echo "Waiting for RustFS..." + sleep 2 + done + /usr/bin/rc ls rustfs/my-bucket >/dev/null 2>&1 || /usr/bin/rc mb rustfs/my-bucket + networks: + - warehouse + + trino: + image: trinodb/trino:435 + volumes: + - ./hive.properties:/etc/trino/catalog/hive.properties:ro + - metastore-data:/data/metastore + depends_on: + create-bucket: + condition: service_completed_successfully + networks: + - warehouse + +networks: + warehouse: + +volumes: + rustfs-data: + metastore-data: +``` + +## 2. Start the deployment + +Resolve the Compose file before starting containers: + +```bash +docker compose config +``` + +Start the services and wait for the bucket initializer to finish: + +```bash +docker compose up -d +docker compose ps -a +``` + +Trino is up when the server log reports `SERVER STARTED`. The container runs as user `trino` (uid 1000); make sure the metastore volume is writable: + +```bash +docker compose exec trino id +docker compose exec trino ls -la /data/metastore +``` + +## 3. Create a schema and a table + +Create the schema without an explicit location — Trino places it under the catalog directory in RustFS: + +```bash +docker compose exec trino trino --execute \ + "CREATE SCHEMA hive.demo" +``` + +Create a table and insert five rows: + +```bash +docker compose exec trino trino --execute \ + "CREATE TABLE hive.demo.events (id bigint, label varchar) WITH (format = 'parquet')" + +docker compose exec trino trino --execute \ + "INSERT INTO hive.demo.events VALUES (1,'alpha'),(2,'bravo'),(3,'charlie'),(4,'delta'),(5,'echo')" +``` + +```text +INSERT: 5 rows +``` + +## 4. Query the data + +Read the rows back: + +```bash +docker compose exec trino trino --execute \ + "SELECT * FROM hive.demo.events ORDER BY id" +``` + +```text +"1","alpha" +"2","bravo" +"3","charlie" +"4","delta" +"5","echo" +``` + +## 5. Verify objects in RustFS + +List the metastore prefix through the bucket-initializer image: + +```bash +docker compose run --rm --entrypoint /bin/sh create-bucket -c \ + '/usr/bin/rc alias set rustfs http://rustfs:9000 "$RUSTFS_ACCESS_KEY" "$RUSTFS_SECRET_KEY" >/dev/null && /usr/bin/rc ls rustfs/my-bucket/trino-metastore --recursive' +``` + +```text +[2026-09-21 01:55:16] 155 B trino-metastore/.demo.trinoSchema +[2026-09-21 01:55:19] 474 B trino-metastore/demo/events/.trinoPermissions/user_trino +[2026-09-21 01:55:25] 1007 B trino-metastore/demo/events/.trinoSchema +[2026-09-21 01:55:25] 432 B trino-metastore/demo/events/20260921_..._cb761cec-...parquet +``` + +You can also browse the prefix in the RustFS Console: + +![The Trino metadata and data objects stored in the RustFS Console](./images/rustfs-trino-objects.png) + +## 6. Use RustFS S3 Tables + +RustFS S3 Tables provides a built-in Apache Iceberg REST catalog, so Trino can treat a table bucket as a managed Iceberg warehouse while the data stays in RustFS. Enable a table bucket and connect Trino's Iceberg connector to the REST catalog as described in [S3 Tables](/administration/data/s3-tables): the REST catalog URI is `http://:9000/iceberg`, the warehouse is the bucket name, and both catalog requests (AWS Signature Version 4, signing name `s3`) and S3 file access use path-style addressing. + +Per the S3 Tables support matrix, Trino has been probed for read-only access against the catalog; validate write compatibility and the exact Trino version you deploy before adopting this path in production. + +## 7. Stop or reset the stack + +Stop the containers while keeping the RustFS data volume: + +```bash +docker compose down +``` + +To delete the stored metadata and data and start from an empty RustFS volume, explicitly include `--volumes`: + +```bash +docker compose down --volumes +``` + +## Troubleshooting + +### Configuration errors for `fs.native-s3.enabled` or `fs.s3.enabled` + +The native S3 filesystem property changed across Trino versions: Trino 435 uses `fs.native-s3.enabled`, newer releases use `fs.s3.enabled`. This guide pins `trinodb/trino:435`, so use `fs.native-s3.enabled`. + +### "Table directory must be ..." when creating a table + +With the file metastore, table locations must stay under `hive.metastore.catalog.dir`. Create the schema without an explicit location, or point the schema location at a directory inside the same bucket prefix. + +### Hive CSV storage format only supports VARCHAR + +The CSV format rejects non-string columns. Use `format = 'parquet'` (as in this guide) for typed tables. + +### AccessDenied or 403 responses + +Confirm that the credentials in `hive.properties` match the RustFS credentials and that the `create-bucket` service completed successfully: + +```bash +docker compose logs create-bucket +``` + +## Next steps + +- Review [S3 compatibility notes](/administration/protocols/s3) before adopting additional S3 operations. +- Create dedicated production credentials with [Access Key Management](/security-compliance/iam/access-token). +- Follow the [Trino documentation](https://trino.io/docs/current/) to connect BI tools and add object storage catalogs. diff --git a/content/fr/developer/integration/big-data/flink.md b/content/fr/developer/integration/big-data/flink.md new file mode 100644 index 00000000..998bd66a --- /dev/null +++ b/content/fr/developer/integration/big-data/flink.md @@ -0,0 +1,278 @@ +--- +title: "Apache Flink" +description: "Lisez et écrivez des données CSV stockées dans le stockage objet RustFS avec Apache Flink et son plugin S3 filesystem." +--- + +Ce guide connecte [Apache Flink](https://github.com/apache/flink) à **RustFS** via le plugin S3 filesystem de Flink (`flink-s3-fs-hadoop`). Vous allez démarrer un cluster de session avec Docker Compose, écrire un résultat borné dans le bucket en mode batch, puis le relire via Flink SQL. Le flux a été validé avec `flink:1.20` et `rustfs/rustfs-x86-musl:v2.3.1`. + +Vous avez besoin de Docker avec le plugin Compose. Ce déploiement est destiné aux tests d'intégration locaux, pas à la production. + +## Architecture + +```mermaid +flowchart LR + Job["Flink SQL job"] -->|"filesystem connector"| S3["S3 plugin (flink-s3-fs-hadoop)"] + S3 -->|"GET / PUT"| RustFS["RustFS :9000"] + Init["init-bucket job"] -->|"create my-bucket"| RustFS +``` + +Le plugin `flink-s3-fs-hadoop` enregistre le schéma `s3://` pour le connecteur filesystem de Flink. Le point de terminaison, l'adressage path-style, HTTP simple et les identifiants sont configurés via les propriétés `s3.*` de `flink-conf.yaml` (passées par `FLINK_PROPERTIES`). + +## 1. Créer les fichiers du projet + +Créez un répertoire de travail : + +```bash +mkdir rustfs-flink +cd rustfs-flink +``` + +Créez un fichier d'environnement et remplacez les deux espaces réservés d'identifiants : + +```ini title=".env" +RUSTFS_ACCESS_KEY= +RUSTFS_SECRET_KEY= +``` + +Utilisez des identifiants dédiés pour le bucket. Ne commettez pas `.env` dans le contrôle de version. + +Le plugin S3 est embarqué dans l'image sous `/opt/flink/opt/` et doit être copié vers `/opt/flink/plugins/s3fs/` pour être chargé. Préparez un répertoire local : + +```bash +mkdir -p s3fs +docker create --name flink-tmp flink:1.20 +docker cp flink-tmp:/opt/flink/opt/flink-s3-fs-hadoop-1.20.5.jar s3fs/ +docker rm flink-tmp +``` + +Créez le fichier Compose : + +```yaml title="compose.yaml" +services: + rustfs: + image: rustfs/rustfs-x86-musl:v2.3.1 + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + RUSTFS_VOLUMES: /data + RUSTFS_ADDRESS: ":9000" + RUSTFS_CONSOLE_ADDRESS: ":9001" + RUSTFS_CONSOLE_ENABLE: "true" + volumes: + - rustfs-data:/data + ports: + - "9000:9000" + - "9001:9001" + healthcheck: + test: ["CMD", "curl", "-sf", "http://127.0.0.1:9000/health"] + interval: 10s + timeout: 5s + retries: 6 + start_period: 10s + networks: + - flink + + create-bucket: + image: rustfs/rc:latest + depends_on: + rustfs: + condition: service_healthy + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + entrypoint: + - /bin/sh + - -c + - | + until /usr/bin/rc alias set rustfs http://rustfs:9000 "$${RUSTFS_ACCESS_KEY}" "$${RUSTFS_SECRET_KEY}"; do + echo "Waiting for RustFS..." + sleep 2 + done + /usr/bin/rc ls rustfs/my-bucket >/dev/null 2>&1 || /usr/bin/rc mb rustfs/my-bucket + networks: + - flink + + jobmanager: + image: flink:1.20 + command: jobmanager + environment: + FLINK_PROPERTIES: | + jobmanager.rpc.address: jobmanager + rest.address: jobmanager + rest.bind-address: 0.0.0.0 + s3.access-key: ${RUSTFS_ACCESS_KEY} + s3.secret-key: ${RUSTFS_SECRET_KEY} + s3.endpoint: http://rustfs:9000 + s3.path-style-access: true + volumes: + - ./s3fs:/opt/flink/plugins/s3fs:ro + ports: + - "8081:8081" + depends_on: + create-bucket: + condition: service_completed_successfully + networks: + - flink + + taskmanager: + image: flink:1.20 + command: taskmanager + environment: + FLINK_PROPERTIES: | + jobmanager.rpc.address: jobmanager + taskmanager.host: taskmanager + s3.access-key: ${RUSTFS_ACCESS_KEY} + s3.secret-key: ${RUSTFS_SECRET_KEY} + s3.endpoint: http://rustfs:9000 + s3.path-style-access: true + volumes: + - ./s3fs:/opt/flink/plugins/s3fs:ro + depends_on: + jobmanager: + condition: service_started + networks: + - flink + +networks: + flink: + +volumes: + rustfs-data: +``` + +Les propriétés `s3.access-key`, `s3.secret-key`, `s3.endpoint` et `s3.path-style-access` configurent le plugin S3 sur le JobManager et le TaskManager. + +## 2. Démarrer le déploiement + +Vérifiez le fichier Compose avant de démarrer les conteneurs : + +```bash +docker compose config +``` + +Démarrez les services et attendez la fin de l'initialisation du bucket : + +```bash +docker compose up -d +docker compose ps -a +``` + +## 3. Écrire un résultat dans RustFS + +Créez le job SQL — mode batch avec le sink filesystem : + +```yaml title="batch.sql" +SET 'execution.runtime-mode' = 'batch'; + +CREATE TABLE sink ( + id INT, + payload STRING +) WITH ( + 'connector' = 'filesystem', + 'path' = 's3://my-bucket/flink-out/', + 'format' = 'csv' +); + +INSERT INTO sink + VALUES (1, 'alpha'), (2, 'bravo'), (3, 'charlie'), (4, 'delta'), (5, 'echo'); +``` + +Soumettez-le via le client SQL dans le JobManager : + +```bash +docker compose exec jobmanager bash -c "/opt/flink/bin/sql-client.sh embedded -f /dev/stdin" < batch.sql +``` + +Le job se termine une fois toutes les lignes écrites. + +## 4. Relire les données + +Créez la requête de lecture — le connecteur filesystem scanne le préfixe : + +```yaml title="read.sql" +CREATE TABLE readings ( + id INT, + payload STRING +) WITH ( + 'connector' = 'filesystem', + 'path' = 's3://my-bucket/flink-out/', + 'format' = 'csv' +); + +SET 'sql-client.execution.result-mode' = 'TABLEAU'; + +SELECT * FROM readings; +``` + +```bash +docker compose exec jobmanager bash -c "/opt/flink/bin/sql-client.sh embedded -f /dev/stdin" < read.sql +``` + +```text ++----+-------------+--------------------------------+ +| op | id | payload | ++----+-------------+--------------------------------+ +| +I | 1 | alpha | +| +I | 2 | bravo | +| +I | 3 | charlie | +| +I | 4 | delta | +| +I | 5 | echo | ++----+-------------+--------------------------------+ +``` + +## 5. Vérifier les objets dans RustFS + +Listez le préfixe via l'image d'initialisation du bucket : + +```bash +docker compose run --rm --entrypoint /bin/sh create-bucket -c \ + '/usr/bin/rc alias set rustfs http://rustfs:9000 "$RUSTFS_ACCESS_KEY" "$RUSTFS_SECRET_KEY" >/dev/null && /usr/bin/rc ls rustfs/my-bucket/flink-out --recursive' +``` + +```text +[2026-09-21 01:34:52] 41 B flink-out/part-f759f9e8-3d1b-46a1-a92e-53e9b727e831-task-0-file-0 +``` + +Vous pouvez également parcourir le préfixe dans la console RustFS : + +![Le fichier de sortie Flink stocké dans la console RustFS](./images/rustfs-flink-objects.png) + +## 6. Utiliser RustFS S3 Tables + +RustFS S3 Tables fournit un catalogue REST Apache Iceberg intégré, permettant à Flink de traiter un table bucket comme un entrepôt Iceberg géré tandis que les données restent dans RustFS. Activez un table bucket et orientez le catalogue REST du connecteur Iceberg de Flink vers RustFS comme décrit dans [S3 Tables](/administration/data/s3-tables) : l'URI du catalogue REST est `http://:9000/iceberg`, le warehouse est le nom du bucket, et les requêtes de catalogue (AWS Signature Version 4, nom de signature `s3`) comme l'accès aux fichiers S3 utilisent l'adressage path-style. + +Selon la matrice de support S3 Tables, validez les versions exactes de Flink et d'Iceberg que vous déployez par rapport au catalogue avant d'adopter ce chemin en production. + +## 7. Arrêter ou réinitialiser la pile + +Arrêtez les conteneurs en conservant le volume de données RustFS : + +```bash +docker compose down +``` + +Pour supprimer les fichiers stockés et repartir d'un volume RustFS vide, ajoutez explicitement `--volumes` : + +```bash +docker compose down --volumes +``` + +## Dépannage + +### No AWS Credentials provided / AccessDenied lors des écritures + +Le plugin S3 lit ses identifiants depuis les propriétés `s3.*` de `flink-conf.yaml`. Vérifiez que `s3.access-key`, `s3.secret-key`, `s3.endpoint` et `s3.path-style-access` sont présents dans `FLINK_PROPERTIES` pour **le** JobManager **et le** TaskManager, et que le jar du plugin existe dans `/opt/flink/plugins/s3fs/` sur chacun. + +### Le TaskManager ne résout pas le nom `rustfs` + +Tous les conteneurs Flink et RustFS doivent partager un réseau Compose. Si RustFS est attaché à un réseau externe, connectez également les conteneurs Flink à ce réseau avant de soumettre le job. + +### Une écriture streaming échoue avec "Stream closed" après un échec + +La reprise d'un upload S3 en cours après un échec peut laisser le writer dans un état irrécupérable. Supprimez le préfixe de sortie du job dans le bucket et soumettez à nouveau, ou utilisez le mode batch comme dans ce guide pour les écritures ponctuelles. + +## Prochaines étapes + +- Consultez les [notes de compatibilité S3](/administration/protocols/s3) avant d'adopter d'autres opérations S3. +- Créez des identifiants de production dédiés avec la [gestion des clés d'accès](/security-compliance/iam/access-token). +- Suivez la [documentation Apache Flink](https://nightlies.apache.org/flink/flink-docs-stable/) pour les options du connecteur filesystem telles que le partitionnement et la compaction. diff --git a/content/fr/developer/integration/big-data/images/rustfs-flink-objects.png b/content/fr/developer/integration/big-data/images/rustfs-flink-objects.png new file mode 100644 index 00000000..caf4f204 Binary files /dev/null and b/content/fr/developer/integration/big-data/images/rustfs-flink-objects.png differ diff --git a/content/fr/developer/integration/big-data/images/rustfs-spark-objects.png b/content/fr/developer/integration/big-data/images/rustfs-spark-objects.png new file mode 100644 index 00000000..f43539d7 Binary files /dev/null and b/content/fr/developer/integration/big-data/images/rustfs-spark-objects.png differ diff --git a/content/fr/developer/integration/big-data/images/rustfs-trino-objects.png b/content/fr/developer/integration/big-data/images/rustfs-trino-objects.png new file mode 100644 index 00000000..3846a0fd Binary files /dev/null and b/content/fr/developer/integration/big-data/images/rustfs-trino-objects.png differ diff --git a/content/fr/developer/integration/big-data/index.md b/content/fr/developer/integration/big-data/index.md index 05b8ae62..11b4660d 100644 --- a/content/fr/developer/integration/big-data/index.md +++ b/content/fr/developer/integration/big-data/index.md @@ -12,5 +12,8 @@ Use **RustFS** as the object storage layer for data analytics systems that suppo - [Milvus](./milvus.md) - [DuckDB](./duckdb.md) - [InfluxDB](./influxdb.md) +- [Spark](./spark.md) +- [Flink](./flink.md) +- [Trino](./trino.md) Keep application data in a dedicated bucket and prefix, and use credentials scoped to the required bucket operations. \ No newline at end of file diff --git a/content/fr/developer/integration/big-data/meta.json b/content/fr/developer/integration/big-data/meta.json index 0b73dddb..1d6f178f 100644 --- a/content/fr/developer/integration/big-data/meta.json +++ b/content/fr/developer/integration/big-data/meta.json @@ -1,10 +1,13 @@ { - "title": "Analyse de données", + "title": "Data Analytics", "pages": [ "iceberg", "pyiceberg", "milvus", "duckdb", - "influxdb" + "influxdb", + "spark", + "flink", + "trino" ] } diff --git a/content/fr/developer/integration/big-data/spark.md b/content/fr/developer/integration/big-data/spark.md new file mode 100644 index 00000000..ff275aab --- /dev/null +++ b/content/fr/developer/integration/big-data/spark.md @@ -0,0 +1,220 @@ +--- +title: "Apache Spark" +description: "Lisez et écrivez des données Parquet stockées dans le stockage objet RustFS avec Apache Spark via le connecteur s3a." +--- + +Ce guide connecte [Apache Spark](https://github.com/apache/spark) à **RustFS** via le connecteur `s3a`. Vous allez démarrer RustFS avec Docker Compose, exécuter un job Spark qui écrit un dataset Parquet dans le bucket, le relire, puis vérifier les objets dans RustFS. Le flux a été validé avec `apache/spark:3.5.6` (Hadoop 3.3.4 via `hadoop-aws`) et `rustfs/rustfs-x86-musl:v2.3.1`. + +Vous avez besoin de Docker avec le plugin Compose. Ce déploiement est destiné aux tests d'intégration locaux, pas à la production. + +## Architecture + +```mermaid +flowchart LR + Spark["Spark driver + executors"] -->|"S3AFileSystem"| RustFS["RustFS :9000"] + Init["init-bucket job"] -->|"create my-bucket"| RustFS +``` + +Spark communique avec RustFS via le système de fichiers S3A de `hadoop-aws`. Les réglages du connecteur — point de terminaison, adressage path-style, HTTP simple et identifiants — sont passés en propriétés `spark.hadoop.fs.s3a.*`. + +## 1. Créer les fichiers du projet + +Créez un répertoire de travail : + +```bash +mkdir rustfs-spark +cd rustfs-spark +``` + +Créez un fichier d'environnement et remplacez les deux espaces réservés d'identifiants : + +```ini title=".env" +RUSTFS_ACCESS_KEY= +RUSTFS_SECRET_KEY= +``` + +Utilisez des identifiants dédiés pour le bucket. Ne commettez pas `.env` dans le contrôle de version. + +Créez le job Spark : + +```python title="job.py" +from pyspark.sql import SparkSession + +spark = SparkSession.builder.appName("rustfs-spark-demo").getOrCreate() +spark.sparkContext.setLogLevel("WARN") + +spark.range(1000).withColumnRenamed("id", "num") \ + .write.mode("overwrite").parquet("s3a://my-bucket/spark-demo/events") + +back = spark.read.parquet("s3a://my-bucket/spark-demo/events") +print("ROWS_READ_BACK:", back.count()) +spark.stop() +``` + +Créez le fichier Compose : + +```yaml title="compose.yaml" +services: + rustfs: + image: rustfs/rustfs-x86-musl:v2.3.1 + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + RUSTFS_VOLUMES: /data + RUSTFS_ADDRESS: ":9000" + RUSTFS_CONSOLE_ADDRESS: ":9001" + RUSTFS_CONSOLE_ENABLE: "true" + volumes: + - rustfs-data:/data + ports: + - "9000:9000" + - "9001:9001" + healthcheck: + test: ["CMD", "curl", "-sf", "http://127.0.0.1:9000/health"] + interval: 10s + timeout: 5s + retries: 6 + start_period: 10s + networks: + - warehouse + + create-bucket: + image: rustfs/rc:latest + depends_on: + rustfs: + condition: service_healthy + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + entrypoint: + - /bin/sh + - -c + - | + until /usr/bin/rc alias set rustfs http://rustfs:9000 "$${RUSTFS_ACCESS_KEY}" "$${RUSTFS_SECRET_KEY}"; do + echo "Waiting for RustFS..." + sleep 2 + done + /usr/bin/rc ls rustfs/my-bucket >/dev/null 2>&1 || /usr/bin/rc mb rustfs/my-bucket + networks: + - warehouse + + spark: + image: apache/spark:3.5.6 + entrypoint: ["/opt/spark/bin/spark-submit"] + volumes: + - ./job.py:/job.py:ro + command: + - --conf + - spark.jars.ivy=/tmp/.ivy2 + - --packages + - org.apache.hadoop:hadoop-aws:3.3.4 + - --conf + - spark.hadoop.fs.s3a.endpoint=http://rustfs:9000 + - --conf + - spark.hadoop.fs.s3a.access.key=${RUSTFS_ACCESS_KEY} + - --conf + - spark.hadoop.fs.s3a.secret.key=${RUSTFS_SECRET_KEY} + - --conf + - spark.hadoop.fs.s3a.path.style.access=true + - --conf + - spark.hadoop.fs.s3a.connection.ssl.enabled=false + - --conf + - spark.hadoop.fs.s3a.impl=org.apache.hadoop.fs.s3a.S3AFileSystem + - /job.py + depends_on: + create-bucket: + condition: service_completed_successfully + networks: + - warehouse + +networks: + warehouse: + +volumes: + rustfs-data: +``` + +`--packages org.apache.hadoop:hadoop-aws:3.3.4` télécharge le connecteur S3 au lancement ; il doit correspondre à la version de Hadoop embarquée dans l'image Spark. `spark.jars.ivy=/tmp/.ivy2` déplace le cache de téléchargement vers un répertoire inscriptible. + +## 2. Démarrer le stockage et exécuter le job + +Démarrez les services de stockage : + +```bash +docker compose up -d +docker compose ps -a +``` + +Exécutez le job Spark : + +```bash +docker compose run --rm spark +``` + +```text +ROWS_READ_BACK: 1000 +``` + +## 3. Vérifier les objets dans RustFS + +Listez le dataset via l'image d'initialisation du bucket : + +```bash +docker compose run --rm --entrypoint /bin/sh create-bucket -c \ + '/usr/bin/rc alias set rustfs http://rustfs:9000 "$RUSTFS_ACCESS_KEY" "$RUSTFS_SECRET_KEY" >/dev/null && /usr/bin/rc ls rustfs/my-bucket/spark-demo --recursive' +``` + +```text +[2026-09-21 00:58:22] 0 B spark-demo/events/_SUCCESS +[2026-09-21 00:58:21] 1.46 KiB spark-demo/events/part-00000-...-c000.snappy.parquet +[2026-09-21 00:58:21] 1.46 KiB spark-demo/events/part-00001-...-c000.snappy.parquet +[2026-09-21 00:58:21] 1.46 KiB spark-demo/events/part-00002-...-c000.snappy.parquet +``` + +Vous pouvez également parcourir le préfixe dans la console RustFS : + +![Sortie Parquet de Spark stockée dans la console RustFS](./images/rustfs-spark-objects.png) + +## 4. Utiliser RustFS S3 Tables + +RustFS S3 Tables fournit un catalogue REST Apache Iceberg intégré, permettant à Spark de traiter un table bucket comme un entrepôt Iceberg géré tandis que les données restent dans RustFS. Activez un table bucket et connectez le catalogue REST Iceberg de Spark comme décrit dans [S3 Tables](/administration/data/s3-tables) : l'URI du catalogue REST est `http://:9000/iceberg`, le warehouse est le nom du bucket, et les requêtes de catalogue (AWS Signature Version 4, nom de signature `s3`) comme l'accès aux fichiers S3 utilisent l'adressage path-style. + +Selon la matrice de support S3 Tables, validez les versions exactes de Spark et d'Iceberg que vous déployez par rapport au catalogue avant d'adopter ce chemin en production. + +## 5. Arrêter ou réinitialiser la pile + +Arrêtez les conteneurs en conservant le volume de données RustFS : + +```bash +docker compose down +``` + +Pour supprimer le dataset et repartir d'un volume RustFS vide, ajoutez explicitement `--volumes` : + +```bash +docker compose down --volumes +``` + +## Dépannage + +### NumberFormatException: For input string: "60s" + +La version de `hadoop-aws` ne correspond pas à la version de Hadoop embarquée dans l'image Spark. Les images Spark 4.x nécessitent `hadoop-aws` 3.4.x ; ce guide épingle `apache/spark:3.5.6` avec `hadoop-aws:3.3.4`. + +### Échecs de connexion à `rustfs:9000` + +`fs.s3a.endpoint` est résolu à l'intérieur du réseau Compose. Pour un processus Spark exécuté sur l'hôte, utilisez `http://localhost:9000`. + +### Réponses AccessDenied ou 403 + +Vérifiez que les réglages du connecteur correspondent aux identifiants RustFS et que la tâche `create-bucket` s'est terminée avec succès : + +```bash +docker compose logs create-bucket +``` + +## Prochaines étapes + +- Consultez les [notes de compatibilité S3](/administration/protocols/s3) avant d'adopter d'autres opérations S3. +- Créez des identifiants de production dédiés avec la [gestion des clés d'accès](/security-compliance/iam/access-token). +- Suivez la [documentation Spark](https://spark.apache.org/docs/latest/) pour le structured streaming et les options de sources de données. diff --git a/content/fr/developer/integration/big-data/trino.md b/content/fr/developer/integration/big-data/trino.md new file mode 100644 index 00000000..9e04da2a --- /dev/null +++ b/content/fr/developer/integration/big-data/trino.md @@ -0,0 +1,250 @@ +--- +title: "Trino" +description: "Interrogez des données CSV et Parquet stockées dans le stockage objet RustFS avec Trino et le file metastore du connecteur hive." +--- + +Ce guide connecte [Trino](https://github.com/trinodb/trino) — le moteur de requêtes SQL distribué — à **RustFS** via le connecteur hive avec son metastore basé sur les fichiers et le système de fichiers S3 natif. Vous allez créer un schéma et une table, insérer des lignes, les relire, puis vérifier les objets dans RustFS. Les métadonnées et les fichiers de données résident tous deux dans RustFS. Le flux a été validé avec `trinodb/trino:435` et `rustfs/rustfs-x86-musl:v2.3.1`. + +Vous avez besoin de Docker avec le plugin Compose. Ce déploiement est destiné aux tests d'intégration locaux, pas à la production. + +## Architecture + +```mermaid +flowchart LR + Client["trino CLI"] -->|"SQL"| Trino["Trino :8080"] + Trino -->|"metadata JSON"| RustFS["RustFS :9000"] + Trino -->|"data files"| RustFS + Init["init-bucket job"] -->|"create my-bucket"| RustFS +``` + +Le connecteur hive avec `hive.metastore=file` conserve les métadonnées de schémas et de tables sous forme d'objets JSON sous le répertoire du catalogue, et le système de fichiers S3 natif (`fs.s3.enabled`) stocke les métadonnées et les données dans RustFS avec un adressage path-style en HTTP simple. + +## 1. Créer les fichiers du projet + +Créez un répertoire de travail : + +```bash +mkdir rustfs-trino +cd rustfs-trino +``` + +Créez un fichier d'environnement et remplacez les deux espaces réservés d'identifiants : + +```ini title=".env" +RUSTFS_ACCESS_KEY= +RUSTFS_SECRET_KEY= +``` + +Utilisez des identifiants dédiés pour le bucket. Ne commettez pas `.env` dans le contrôle de version. + +Créez la configuration du catalogue pour Trino : + +```ini title="hive.properties" +connector.name=hive +hive.metastore=file +hive.metastore.catalog.dir=s3://my-bucket/trino-metastore +fs.s3.enabled=true +s3.endpoint=http://rustfs:9000 +s3.region=us-east-1 +s3.path-style-access=true +s3.aws-access-key= +s3.aws-secret-key= +``` + +`hive.metastore.catalog.dir` pointe le file metastore dans le bucket : métadonnées et données résident donc dans RustFS. `fs.s3.enabled` active le système de fichiers S3 natif ; `s3.path-style-access` est requis pour le point de terminaison du réseau de conteneurs. + +Créez le fichier Compose : + +```yaml title="compose.yaml" +services: + rustfs: + image: rustfs/rustfs-x86-musl:v2.3.1 + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + RUSTFS_VOLUMES: /data + RUSTFS_ADDRESS: ":9000" + RUSTFS_CONSOLE_ADDRESS: ":9001" + RUSTFS_CONSOLE_ENABLE: "true" + volumes: + - rustfs-data:/data + ports: + - "9000:9000" + - "9001:9001" + healthcheck: + test: ["CMD", "curl", "-sf", "http://127.0.0.1:9000/health"] + interval: 10s + timeout: 5s + retries: 6 + start_period: 10s + networks: + - warehouse + + create-bucket: + image: rustfs/rc:latest + depends_on: + rustfs: + condition: service_healthy + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + entrypoint: + - /bin/sh + - -c + - | + until /usr/bin/rc alias set rustfs http://rustfs:9000 "$${RUSTFS_ACCESS_KEY}" "$${RUSTFS_SECRET_KEY}"; do + echo "Waiting for RustFS..." + sleep 2 + done + /usr/bin/rc ls rustfs/my-bucket >/dev/null 2>&1 || /usr/bin/rc mb rustfs/my-bucket + networks: + - warehouse + + trino: + image: trinodb/trino:435 + volumes: + - ./hive.properties:/etc/trino/catalog/hive.properties:ro + - metastore-data:/data/metastore + depends_on: + create-bucket: + condition: service_completed_successfully + networks: + - warehouse + +networks: + warehouse: + +volumes: + rustfs-data: + metastore-data: +``` + +## 2. Démarrer le déploiement + +Vérifiez le fichier Compose avant de démarrer les conteneurs : + +```bash +docker compose config +``` + +Démarrez les services et attendez la fin de l'initialisation du bucket : + +```bash +docker compose up -d +docker compose ps -a +``` + +Trino est démarré quand le journal du serveur indique `SERVER STARTED`. Le conteneur s'exécute sous l'utilisateur `trino` (uid 1000) ; assurez-vous que le volume du metastore est inscriptible : + +```bash +docker compose exec trino id +docker compose exec trino ls -la /data/metastore +``` + +## 3. Créer un schéma et une table + +Créez le schéma sans location explicite — Trino le place sous le répertoire du catalogue dans RustFS : + +```bash +docker compose exec trino trino --execute \ + "CREATE SCHEMA hive.demo" +``` + +Créez une table et insérez cinq lignes : + +```bash +docker compose exec trino trino --execute \ + "CREATE TABLE hive.demo.events (id bigint, label varchar) WITH (format = 'parquet')" + +docker compose exec trino trino --execute \ + "INSERT INTO hive.demo.events VALUES (1,'alpha'),(2,'bravo'),(3,'charlie'),(4,'delta'),(5,'echo')" +``` + +```text +INSERT: 5 rows +``` + +## 4. Interroger les données + +Relisez les lignes : + +```bash +docker compose exec trino trino --execute \ + "SELECT * FROM hive.demo.events ORDER BY id" +``` + +```text +"1","alpha" +"2","bravo" +"3","charlie" +"4","delta" +"5","echo" +``` + +## 5. Vérifier les objets dans RustFS + +Listez le préfixe du metastore via l'image d'initialisation du bucket : + +```bash +docker compose run --rm --entrypoint /bin/sh create-bucket -c \ + '/usr/bin/rc alias set rustfs http://rustfs:9000 "$RUSTFS_ACCESS_KEY" "$RUSTFS_SECRET_KEY" >/dev/null && /usr/bin/rc ls rustfs/my-bucket/trino-metastore --recursive' +``` + +```text +[2026-09-21 01:55:16] 155 B trino-metastore/.demo.trinoSchema +[2026-09-21 01:55:19] 474 B trino-metastore/demo/events/.trinoPermissions/user_trino +[2026-09-21 01:55:25] 1007 B trino-metastore/demo/events/.trinoSchema +[2026-09-21 01:55:25] 432 B trino-metastore/demo/events/20260921_..._cb761cec-...parquet +``` + +Vous pouvez également parcourir le préfixe dans la console RustFS : + +![Les objets de métadonnées et de données Trino stockés dans la console RustFS](./images/rustfs-trino-objects.png) + +## 6. Utiliser RustFS S3 Tables + +RustFS S3 Tables fournit un catalogue REST Apache Iceberg intégré, permettant à Trino de traiter un table bucket comme un entrepôt Iceberg géré tandis que les données restent dans RustFS. Activez un table bucket et connectez le connecteur Iceberg de Trino au catalogue REST comme décrit dans [S3 Tables](/administration/data/s3-tables) : l'URI du catalogue REST est `http://:9000/iceberg`, le warehouse est le nom du bucket, et les requêtes de catalogue (AWS Signature Version 4, nom de signature `s3`) comme l'accès aux fichiers S3 utilisent l'adressage path-style. + +Selon la matrice de support S3 Tables, Trino a fait l'objet d'une sonde en lecture seule contre le catalogue ; validez la compatibilité en écriture et la version exacte de Trino que vous déployez avant d'adopter ce chemin en production. + +## 7. Arrêter ou réinitialiser la pile + +Arrêtez les conteneurs en conservant le volume de données RustFS : + +```bash +docker compose down +``` + +Pour supprimer les métadonnées et données stockées et repartir d'un volume RustFS vide, ajoutez explicitement `--volumes` : + +```bash +docker compose down --volumes +``` + +## Dépannage + +### Erreurs de configuration pour `fs.native-s3.enabled` ou `fs.s3.enabled` + +Le nom de la propriété du système de fichiers S3 natif a changé selon les versions de Trino : Trino 435 utilise `fs.native-s3.enabled`, les versions plus récentes `fs.s3.enabled`. Ce guide épingle `trinodb/trino:435`, utilisez donc `fs.native-s3.enabled`. + +### "Table directory must be ..." lors de la création d'une table + +Avec le file metastore, les locations de tables doivent rester sous `hive.metastore.catalog.dir`. Créez le schéma sans location explicite, ou pointez la location du schéma vers un répertoire du même préfixe de bucket. + +### Hive CSV storage format only supports VARCHAR + +Le format CSV rejette les colonnes non textuelles. Utilisez `format = 'parquet'` pour les tables typées, comme dans ce guide. + +### Réponses AccessDenied ou 403 + +Vérifiez que les identifiants de `hive.properties` correspondent aux identifiants RustFS et que la tâche `create-bucket` s'est terminée avec succès : + +```bash +docker compose logs create-bucket +``` + +## Prochaines étapes + +- Consultez les [notes de compatibilité S3](/administration/protocols/s3) avant d'adopter d'autres opérations S3. +- Créez des identifiants de production dédiés avec la [gestion des clés d'accès](/security-compliance/iam/access-token). +- Suivez la [documentation Trino](https://trino.io/docs/current/) pour connecter des outils BI et ajouter des catalogues de stockage objet. diff --git a/content/ja/developer/integration/big-data/flink.md b/content/ja/developer/integration/big-data/flink.md new file mode 100644 index 00000000..bac30db1 --- /dev/null +++ b/content/ja/developer/integration/big-data/flink.md @@ -0,0 +1,278 @@ +--- +title: "Apache Flink" +description: "RustFS オブジェクトストレージ内の CSV データを Apache Flink とその S3 ファイルシステムプラグインで読み書きします。" +--- + +このガイドでは、Flink の S3 ファイルシステムプラグイン(`flink-s3-fs-hadoop`)を通じて [Apache Flink](https://github.com/apache/flink) を **RustFS** に接続します。Docker Compose でセッションクラスタを起動し、バッチモードで有界结果セットをバケットに書き込み、Flink SQL で読み戻します。この流れは `flink:1.20` と `rustfs/rustfs-x86-musl:v2.3.1` で検証済みです。 + +Docker と Compose プラグインが必要です。このデプロイはローカルでの統合テストを目的としており、本番環境向けではありません。 + +## アーキテクチャ + +```mermaid +flowchart LR + Job["Flink SQL job"] -->|"filesystem connector"| S3["S3 plugin (flink-s3-fs-hadoop)"] + S3 -->|"GET / PUT"| RustFS["RustFS :9000"] + Init["init-bucket job"] -->|"create my-bucket"| RustFS +``` + +`flink-s3-fs-hadoop` プラグインは、Flink の filesystem コネクタ向けに `s3://` スキームを登録します。エンドポイント、パススタイルアドレス指定、平文 HTTP、認証情報は、`FLINK_PROPERTIES` 経由で渡す `flink-conf.yaml` の `s3.*` プロパティで設定します。 + +## 1. プロジェクトファイルを作成する + +作業ディレクトリを作成します。 + +```bash +mkdir rustfs-flink +cd rustfs-flink +``` + +環境変数ファイルを作成し、2 つの認証情報プレースホルダーを置き換えます。 + +```ini title=".env" +RUSTFS_ACCESS_KEY= +RUSTFS_SECRET_KEY= +``` + +バケットには専用の認証情報を使用してください。`.env` をバージョン管理にコミットしないでください。 + +S3 プラグインはイメージ内の `/opt/flink/opt/` に同梱されており、ロードするには `/opt/flink/plugins/s3fs/` へコピーする必要があります。ローカルディレクトリを準備します。 + +```bash +mkdir -p s3fs +docker create --name flink-tmp flink:1.20 +docker cp flink-tmp:/opt/flink/opt/flink-s3-fs-hadoop-1.20.5.jar s3fs/ +docker rm flink-tmp +``` + +Compose ファイルを作成します。 + +```yaml title="compose.yaml" +services: + rustfs: + image: rustfs/rustfs-x86-musl:v2.3.1 + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + RUSTFS_VOLUMES: /data + RUSTFS_ADDRESS: ":9000" + RUSTFS_CONSOLE_ADDRESS: ":9001" + RUSTFS_CONSOLE_ENABLE: "true" + volumes: + - rustfs-data:/data + ports: + - "9000:9000" + - "9001:9001" + healthcheck: + test: ["CMD", "curl", "-sf", "http://127.0.0.1:9000/health"] + interval: 10s + timeout: 5s + retries: 6 + start_period: 10s + networks: + - flink + + create-bucket: + image: rustfs/rc:latest + depends_on: + rustfs: + condition: service_healthy + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + entrypoint: + - /bin/sh + - -c + - | + until /usr/bin/rc alias set rustfs http://rustfs:9000 "$${RUSTFS_ACCESS_KEY}" "$${RUSTFS_SECRET_KEY}"; do + echo "Waiting for RustFS..." + sleep 2 + done + /usr/bin/rc ls rustfs/my-bucket >/dev/null 2>&1 || /usr/bin/rc mb rustfs/my-bucket + networks: + - flink + + jobmanager: + image: flink:1.20 + command: jobmanager + environment: + FLINK_PROPERTIES: | + jobmanager.rpc.address: jobmanager + rest.address: jobmanager + rest.bind-address: 0.0.0.0 + s3.access-key: ${RUSTFS_ACCESS_KEY} + s3.secret-key: ${RUSTFS_SECRET_KEY} + s3.endpoint: http://rustfs:9000 + s3.path-style-access: true + volumes: + - ./s3fs:/opt/flink/plugins/s3fs:ro + ports: + - "8081:8081" + depends_on: + create-bucket: + condition: service_completed_successfully + networks: + - flink + + taskmanager: + image: flink:1.20 + command: taskmanager + environment: + FLINK_PROPERTIES: | + jobmanager.rpc.address: jobmanager + taskmanager.host: taskmanager + s3.access-key: ${RUSTFS_ACCESS_KEY} + s3.secret-key: ${RUSTFS_SECRET_KEY} + s3.endpoint: http://rustfs:9000 + s3.path-style-access: true + volumes: + - ./s3fs:/opt/flink/plugins/s3fs:ro + depends_on: + jobmanager: + condition: service_started + networks: + - flink + +networks: + flink: + +volumes: + rustfs-data: +``` + +`s3.access-key`、`s3.secret-key`、`s3.endpoint`、`s3.path-style-access` プロパティが、JobManager と TaskManager の両方の S3 プラグインを設定します。 + +## 2. デプロイを起動する + +コンテナを起動する前に Compose ファイルを検証します。 + +```bash +docker compose config +``` + +サービスを起動し、バケット初期化の完了を待ちます。 + +```bash +docker compose up -d +docker compose ps -a +``` + +## 3. 结果セットを RustFS に書き込む + +SQL ジョブを作成します。バッチモードと filesystem シンクを使います。 + +```yaml title="batch.sql" +SET 'execution.runtime-mode' = 'batch'; + +CREATE TABLE sink ( + id INT, + payload STRING +) WITH ( + 'connector' = 'filesystem', + 'path' = 's3://my-bucket/flink-out/', + 'format' = 'csv' +); + +INSERT INTO sink + VALUES (1, 'alpha'), (2, 'bravo'), (3, 'charlie'), (4, 'delta'), (5, 'echo'); +``` + +JobManager 内の SQL クライアントから送信します。 + +```bash +docker compose exec jobmanager bash -c "/opt/flink/bin/sql-client.sh embedded -f /dev/stdin" < batch.sql +``` + +すべての行が書き込まれるとジョブは終了します。 + +## 4. データを読み戻す + +読み取りクエリを作成します。filesystem コネクタがプレフィックスを走査します。 + +```yaml title="read.sql" +CREATE TABLE readings ( + id INT, + payload STRING +) WITH ( + 'connector' = 'filesystem', + 'path' = 's3://my-bucket/flink-out/', + 'format' = 'csv' +); + +SET 'sql-client.execution.result-mode' = 'TABLEAU'; + +SELECT * FROM readings; +``` + +```bash +docker compose exec jobmanager bash -c "/opt/flink/bin/sql-client.sh embedded -f /dev/stdin" < read.sql +``` + +```text ++----+-------------+--------------------------------+ +| op | id | payload | ++----+-------------+--------------------------------+ +| +I | 1 | alpha | +| +I | 2 | bravo | +| +I | 3 | charlie | +| +I | 4 | delta | +| +I | 5 | echo | ++----+-------------+--------------------------------+ +``` + +## 5. RustFS 内のオブジェクトを確認する + +バケット初期化イメージを使ってプレフィックスを一覧表示します。 + +```bash +docker compose run --rm --entrypoint /bin/sh create-bucket -c \ + '/usr/bin/rc alias set rustfs http://rustfs:9000 "$RUSTFS_ACCESS_KEY" "$RUSTFS_SECRET_KEY" >/dev/null && /usr/bin/rc ls rustfs/my-bucket/flink-out --recursive' +``` + +```text +[2026-09-21 01:34:52] 41 B flink-out/part-f759f9e8-3d1b-46a1-a92e-53e9b727e831-task-0-file-0 +``` + +RustFS コンソールでこのプレフィックスを参照することもできます。 + +![RustFS コンソールに保存された Flink の出力ファイル](./images/rustfs-flink-objects.png) + +## 6. RustFS S3 Tables を使用する + +RustFS S3 Tables は組み込みの Apache Iceberg REST カタログを提供します。Flink はテーブルバケットを管理された Iceberg ウェアハウスとして扱え、データは RustFS 内に保持されます。[S3 Tables](/administration/data/s3-tables) の説明に従ってテーブルバケットを有効化し、Flink Iceberg コネクタの REST カタログを RustFS に向けてください。REST カタログ URI は `http://:9000/iceberg`、ウェアハウスはバケット名で、カタログリクエスト(AWS Signature Version 4、署名名 `s3`)と S3 ファイルアクセスの両方がパススタイルのアドレス指定を使います。 + +S3 Tables のサポートマトリクスに従い、本番でこの経路を採用する前に、実際にデプロイする Flink と Iceberg のバージョンでカタログを検証してください。 + +## 7. スタックを停止・リセットする + +RustFS データボリュームを保持したままコンテナを停止します。 + +```bash +docker compose down +``` + +保存したファイルを削除して空の RustFS ボリュームからやり直す場合は、明示的に `--volumes` を付けます。 + +```bash +docker compose down --volumes +``` + +## トラブルシューティング + +### 書き込みで No AWS Credentials provided / AccessDenied が出る + +S3 プラグインは `flink-conf.yaml` の `s3.*` プロパティから認証情報を読み込みます。JobManager と TaskManager **両方**の `FLINK_PROPERTIES` に `s3.access-key`、`s3.secret-key`、`s3.endpoint`、`s3.path-style-access` があること、そして各コンテナの `/opt/flink/plugins/s3fs/` にプラグイン jar があることを確認してください。 + +### TaskManager が `rustfs` ホストを解決できない + +すべての Flink コンテナと RustFS は同じ Compose ネットワークを共有する必要があります。RustFS を外部ネットワークに接続している場合は、ジョブを投入する前に Flink コンテナもそのネットワークに接続してください。 + +### 失敗したストリーミング書き込みが "Stream closed" で復帰しない + +失敗後の進行中 S3 アップロードのリカバリでライターが復帰不能な状態になることがあります。バケット内のジョブ出力プレフィックスを削除してジョブを再投入するか、本ガイドのように一回限りの書き込みにはバッチモードを使用してください。 + +## 次のステップ + +- 追加の S3 オペレーションを採用する前に、[S3 互換性ノート](/administration/protocols/s3)を確認してください。 +- [アクセスキー管理](/security-compliance/iam/access-token)で本番用の専用認証情報を作成してください。 +- [Apache Flink ドキュメント](https://nightlies.apache.org/flink/flink-docs-stable/)で、パーティショニングやコンパクションなど filesystem コネクタのオプションを確認してください。 diff --git a/content/ja/developer/integration/big-data/images/rustfs-flink-objects.png b/content/ja/developer/integration/big-data/images/rustfs-flink-objects.png new file mode 100644 index 00000000..caf4f204 Binary files /dev/null and b/content/ja/developer/integration/big-data/images/rustfs-flink-objects.png differ diff --git a/content/ja/developer/integration/big-data/images/rustfs-spark-objects.png b/content/ja/developer/integration/big-data/images/rustfs-spark-objects.png new file mode 100644 index 00000000..f43539d7 Binary files /dev/null and b/content/ja/developer/integration/big-data/images/rustfs-spark-objects.png differ diff --git a/content/ja/developer/integration/big-data/images/rustfs-trino-objects.png b/content/ja/developer/integration/big-data/images/rustfs-trino-objects.png new file mode 100644 index 00000000..3846a0fd Binary files /dev/null and b/content/ja/developer/integration/big-data/images/rustfs-trino-objects.png differ diff --git a/content/ja/developer/integration/big-data/index.md b/content/ja/developer/integration/big-data/index.md index d0848f75..22510c10 100644 --- a/content/ja/developer/integration/big-data/index.md +++ b/content/ja/developer/integration/big-data/index.md @@ -12,5 +12,8 @@ Use **RustFS** as the object storage layer for data analytics systems that suppo - [Milvus](./milvus.md) - [DuckDB](./duckdb.md) - [InfluxDB](./influxdb.md) +- [Spark](./spark.md) +- [Flink](./flink.md) +- [Trino](./trino.md) Keep application data in a dedicated bucket and prefix, and use credentials scoped to the required bucket operations. \ No newline at end of file diff --git a/content/ja/developer/integration/big-data/meta.json b/content/ja/developer/integration/big-data/meta.json index ff6cc556..b2fbd258 100644 --- a/content/ja/developer/integration/big-data/meta.json +++ b/content/ja/developer/integration/big-data/meta.json @@ -5,6 +5,9 @@ "pyiceberg", "milvus", "duckdb", - "influxdb" + "influxdb", + "spark", + "flink", + "trino" ] } diff --git a/content/ja/developer/integration/big-data/spark.md b/content/ja/developer/integration/big-data/spark.md new file mode 100644 index 00000000..d18dc607 --- /dev/null +++ b/content/ja/developer/integration/big-data/spark.md @@ -0,0 +1,220 @@ +--- +title: "Apache Spark" +description: "s3a コネクタ経由で、RustFS オブジェクトストレージ内の Parquet データを Apache Spark で読み書きします。" +--- + +このガイドでは、`s3a` コネクタを通じて [Apache Spark](https://github.com/apache/spark) を **RustFS** に接続します。Docker Compose で RustFS を起動し、Parquet データセットをバケットに書き込んで読み戻す Spark ジョブを実行し、RustFS 内のオブジェクトを確認します。この流れは `apache/spark:3.5.6`(`hadoop-aws` 3.3.4 を同梱)と `rustfs/rustfs-x86-musl:v2.3.1` で検証済みです。 + +Docker と Compose プラグインが必要です。このデプロイはローカルでの統合テストを目的としており、本番環境向けではありません。 + +## アーキテクチャ + +```mermaid +flowchart LR + Spark["Spark driver + executors"] -->|"S3AFileSystem"| RustFS["RustFS :9000"] + Init["init-bucket job"] -->|"create my-bucket"| RustFS +``` + +Spark は `hadoop-aws` の S3A ファイルシステムを通じて RustFS と通信します。エンドポイント、パススタイルアドレス指定、平文 HTTP、認証情報といったコネクタ設定は `spark.hadoop.fs.s3a.*` プロパティとして渡します。 + +## 1. プロジェクトファイルを作成する + +作業ディレクトリを作成します。 + +```bash +mkdir rustfs-spark +cd rustfs-spark +``` + +環境変数ファイルを作成し、2 つの認証情報プレースホルダーを置き換えます。 + +```ini title=".env" +RUSTFS_ACCESS_KEY= +RUSTFS_SECRET_KEY= +``` + +バケットには専用の認証情報を使用してください。`.env` をバージョン管理にコミットしないでください。 + +Spark ジョブを作成します。 + +```python title="job.py" +from pyspark.sql import SparkSession + +spark = SparkSession.builder.appName("rustfs-spark-demo").getOrCreate() +spark.sparkContext.setLogLevel("WARN") + +spark.range(1000).withColumnRenamed("id", "num") \ + .write.mode("overwrite").parquet("s3a://my-bucket/spark-demo/events") + +back = spark.read.parquet("s3a://my-bucket/spark-demo/events") +print("ROWS_READ_BACK:", back.count()) +spark.stop() +``` + +Compose ファイルを作成します。 + +```yaml title="compose.yaml" +services: + rustfs: + image: rustfs/rustfs-x86-musl:v2.3.1 + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + RUSTFS_VOLUMES: /data + RUSTFS_ADDRESS: ":9000" + RUSTFS_CONSOLE_ADDRESS: ":9001" + RUSTFS_CONSOLE_ENABLE: "true" + volumes: + - rustfs-data:/data + ports: + - "9000:9000" + - "9001:9001" + healthcheck: + test: ["CMD", "curl", "-sf", "http://127.0.0.1:9000/health"] + interval: 10s + timeout: 5s + retries: 6 + start_period: 10s + networks: + - warehouse + + create-bucket: + image: rustfs/rc:latest + depends_on: + rustfs: + condition: service_healthy + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + entrypoint: + - /bin/sh + - -c + - | + until /usr/bin/rc alias set rustfs http://rustfs:9000 "$${RUSTFS_ACCESS_KEY}" "$${RUSTFS_SECRET_KEY}"; do + echo "Waiting for RustFS..." + sleep 2 + done + /usr/bin/rc ls rustfs/my-bucket >/dev/null 2>&1 || /usr/bin/rc mb rustfs/my-bucket + networks: + - warehouse + + spark: + image: apache/spark:3.5.6 + entrypoint: ["/opt/spark/bin/spark-submit"] + volumes: + - ./job.py:/job.py:ro + command: + - --conf + - spark.jars.ivy=/tmp/.ivy2 + - --packages + - org.apache.hadoop:hadoop-aws:3.3.4 + - --conf + - spark.hadoop.fs.s3a.endpoint=http://rustfs:9000 + - --conf + - spark.hadoop.fs.s3a.access.key=${RUSTFS_ACCESS_KEY} + - --conf + - spark.hadoop.fs.s3a.secret.key=${RUSTFS_SECRET_KEY} + - --conf + - spark.hadoop.fs.s3a.path.style.access=true + - --conf + - spark.hadoop.fs.s3a.connection.ssl.enabled=false + - --conf + - spark.hadoop.fs.s3a.impl=org.apache.hadoop.fs.s3a.S3AFileSystem + - /job.py + depends_on: + create-bucket: + condition: service_completed_successfully + networks: + - warehouse + +networks: + warehouse: + +volumes: + rustfs-data: +``` + +`--packages org.apache.hadoop:hadoop-aws:3.3.4` は起動時に S3 コネクタをダウンロードします。Spark イメージに同梱の Hadoop バージョンと一致させる必要があります。`spark.jars.ivy=/tmp/.ivy2` はダウンロードキャッシュを書き込み可能なディレクトリへ移動します。 + +## 2. ストレージを起動してジョブを実行する + +ストレージサービスを起動します。 + +```bash +docker compose up -d +docker compose ps -a +``` + +Spark ジョブを実行します。 + +```bash +docker compose run --rm spark +``` + +```text +ROWS_READ_BACK: 1000 +``` + +## 3. RustFS 内のオブジェクトを確認する + +バケット初期化イメージを使ってデータセットを一覧表示します。 + +```bash +docker compose run --rm --entrypoint /bin/sh create-bucket -c \ + '/usr/bin/rc alias set rustfs http://rustfs:9000 "$RUSTFS_ACCESS_KEY" "$RUSTFS_SECRET_KEY" >/dev/null && /usr/bin/rc ls rustfs/my-bucket/spark-demo --recursive' +``` + +```text +[2026-09-21 00:58:22] 0 B spark-demo/events/_SUCCESS +[2026-09-21 00:58:21] 1.46 KiB spark-demo/events/part-00000-...-c000.snappy.parquet +[2026-09-21 00:58:21] 1.46 KiB spark-demo/events/part-00001-...-c000.snappy.parquet +[2026-09-21 00:58:21] 1.46 KiB spark-demo/events/part-00002-...-c000.snappy.parquet +``` + +RustFS コンソールでこのプレフィックスを参照することもできます。 + +![RustFS コンソールに保存された Spark の Parquet 出力](./images/rustfs-spark-objects.png) + +## 4. RustFS S3 Tables を使用する + +RustFS S3 Tables は組み込みの Apache Iceberg REST カタログを提供します。Flink と同様に、Spark はテーブルバケットを管理された Iceberg ウェアハウスとして扱え、データは RustFS 内に保持されます。[S3 Tables](/administration/data/s3-tables) の説明に従ってテーブルバケットを有効化し、Spark の Iceberg REST カタログを接続してください。REST カタログ URI は `http://:9000/iceberg`、ウェアハウスはバケット名で、カタログリクエスト(AWS Signature Version 4、署名名 `s3`)と S3 ファイルアクセスの両方がパススタイルのアドレス指定を使います。 + +S3 Tables のサポートマトリクスに従い、本番でこの経路を採用する前に、実際にデプロイする Spark と Iceberg のバージョンでカタログを検証してください。 + +## 5. スタックを停止・リセットする + +RustFS データボリュームを保持したままコンテナを停止します。 + +```bash +docker compose down +``` + +データセットを削除して空の RustFS ボリュームからやり直す場合は、明示的に `--volumes` を付けます。 + +```bash +docker compose down --volumes +``` + +## トラブルシューティング + +### NumberFormatException: For input string: "60s" + +`hadoop-aws` のバージョンが Spark イメージ同梱の Hadoop バージョンと一致していません。Spark 4.x イメージには `hadoop-aws` 3.4.x が必要です。このガイドは `apache/spark:3.5.6` と `hadoop-aws:3.3.4` の組み合わせに固定しています。 + +### `rustfs:9000` への接続エラー + +`fs.s3a.endpoint` は Compose ネットワーク内で解決されます。ホスト上で実行する Spark プロセスからは `http://localhost:9000` を使用してください。 + +### AccessDenied や 403 レスポンス + +コネクタ設定が RustFS の認証情報と一致しているか、`create-bucket` ジョブが正常に完了しているかを確認してください。 + +```bash +docker compose logs create-bucket +``` + +## 次のステップ + +- 追加の S3 オペレーションを採用する前に、[S3 互換性ノート](/administration/protocols/s3)を確認してください。 +- [アクセスキー管理](/security-compliance/iam/access-token)で本番用の専用認証情報を作成してください。 +- [Spark ドキュメント](https://spark.apache.org/docs/latest/)で構造化ストリーミングやデータソースのオプションを確認してください。 diff --git a/content/ja/developer/integration/big-data/trino.md b/content/ja/developer/integration/big-data/trino.md new file mode 100644 index 00000000..bed0a998 --- /dev/null +++ b/content/ja/developer/integration/big-data/trino.md @@ -0,0 +1,250 @@ +--- +title: "Trino" +description: "hive コネクタのファイルメタストアとネイティブ S3 ファイルシステムを使って、RustFS オブジェクトストレージ内の CSV・Parquet データを Trino で照会します。" +--- + +このガイドでは、分散 SQL クエリエンジンである [Trino](https://github.com/trinodb/trino) を、hive コネクタのファイルメタストアとネイティブ S3 ファイルシステムを通じて **RustFS** に接続します。スキーマとテーブルを作成し、行を插入して読み戻し、RustFS 内のオブジェクトを確認します。テーブルのメタデータもデータファイルも RustFS 内に保存されます。この流れは `trinodb/trino:435` と `rustfs/rustfs-x86-musl:v2.3.1` で検証済みです。 + +Docker と Compose プラグインが必要です。このデプロイはローカルでの統合テストを目的としており、本番環境向けではありません。 + +## アーキテクチャ + +```mermaid +flowchart LR + Client["trino CLI"] -->|"SQL"| Trino["Trino :8080"] + Trino -->|"metadata JSON"| RustFS["RustFS :9000"] + Trino -->|"data files"| RustFS + Init["init-bucket job"] -->|"create my-bucket"| RustFS +``` + +`hive.metastore=file` の hive コネクタは、スキーマとテーブルのメタデータをカタログディレクトリ配下の JSON オブジェクトとして保持し、ネイティブ S3 ファイルシステム(`fs.s3.enabled`)がメタデータとデータファイルの両方を、平文 HTTP 上のパススタイルアドレス指定で RustFS に保存します。 + +## 1. プロジェクトファイルを作成する + +作業ディレクトリを作成します。 + +```bash +mkdir rustfs-trino +cd rustfs-trino +``` + +環境変数ファイルを作成し、2 つの認証情報プレースホルダーを置き換えます。 + +```ini title=".env" +RUSTFS_ACCESS_KEY= +RUSTFS_SECRET_KEY= +``` + +バケットには専用の認証情報を使用してください。`.env` をバージョン管理にコミットしないでください。 + +Trino の catalog 設定を作成します。 + +```ini title="hive.properties" +connector.name=hive +hive.metastore=file +hive.metastore.catalog.dir=s3://my-bucket/trino-metastore +fs.s3.enabled=true +s3.endpoint=http://rustfs:9000 +s3.region=us-east-1 +s3.path-style-access=true +s3.aws-access-key= +s3.aws-secret-key= +``` + +`hive.metastore.catalog.dir` がファイルメタストアをバケット内に向けます。そのためメタデータもデータも RustFS 内に保存されます。`fs.s3.enabled` がネイティブ S3 ファイルシステムを有効化し、コンテナネットワークのエンドポイントには `s3.path-style-access` が必要です。 + +Compose ファイルを作成します。 + +```yaml title="compose.yaml" +services: + rustfs: + image: rustfs/rustfs-x86-musl:v2.3.1 + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + RUSTFS_VOLUMES: /data + RUSTFS_ADDRESS: ":9000" + RUSTFS_CONSOLE_ADDRESS: ":9001" + RUSTFS_CONSOLE_ENABLE: "true" + volumes: + - rustfs-data:/data + ports: + - "9000:9000" + - "9001:9001" + healthcheck: + test: ["CMD", "curl", "-sf", "http://127.0.0.1:9000/health"] + interval: 10s + timeout: 5s + retries: 6 + start_period: 10s + networks: + - warehouse + + create-bucket: + image: rustfs/rc:latest + depends_on: + rustfs: + condition: service_healthy + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + entrypoint: + - /bin/sh + - -c + - | + until /usr/bin/rc alias set rustfs http://rustfs:9000 "$${RUSTFS_ACCESS_KEY}" "$${RUSTFS_SECRET_KEY}"; do + echo "Waiting for RustFS..." + sleep 2 + done + /usr/bin/rc ls rustfs/my-bucket >/dev/null 2>&1 || /usr/bin/rc mb rustfs/my-bucket + networks: + - warehouse + + trino: + image: trinodb/trino:435 + volumes: + - ./hive.properties:/etc/trino/catalog/hive.properties:ro + - metastore-data:/data/metastore + depends_on: + create-bucket: + condition: service_completed_successfully + networks: + - warehouse + +networks: + warehouse: + +volumes: + rustfs-data: + metastore-data: +``` + +## 2. デプロイを起動する + +コンテナを起動する前に Compose ファイルを検証します。 + +```bash +docker compose config +``` + +サービスを起動し、バケット初期化の完了を待ちます。 + +```bash +docker compose up -d +docker compose ps -a +``` + +Trino はサーバーログに `SERVER STARTED` と出力されると起動完了です。コンテナは `trino` ユーザー(uid 1000)で実行されるため、メタストアボリュームが書き込み可能であることを確認してください。 + +```bash +docker compose exec trino id +docker compose exec trino ls -la /data/metastore +``` + +## 3. スキーマとテーブルを作成する + +明示的なロケーションを指定せずにスキーマを作成します。Trino が RustFS 内のカタログディレクトリ配下に配置します。 + +```bash +docker compose exec trino trino --execute \ + "CREATE SCHEMA hive.demo" +``` + +テーブルを作成して 5 行を插入します。 + +```bash +docker compose exec trino trino --execute \ + "CREATE TABLE hive.demo.events (id bigint, label varchar) WITH (format = 'parquet')" + +docker compose exec trino trino --execute \ + "INSERT INTO hive.demo.events VALUES (1,'alpha'),(2,'bravo'),(3,'charlie'),(4,'delta'),(5,'echo')" +``` + +```text +INSERT: 5 rows +``` + +## 4. データを照会する + +行を読み戻します。 + +```bash +docker compose exec trino trino --execute \ + "SELECT * FROM hive.demo.events ORDER BY id" +``` + +```text +"1","alpha" +"2","bravo" +"3","charlie" +"4","delta" +"5","echo" +``` + +## 5. RustFS 内のオブジェクトを確認する + +バケット初期化イメージを使ってメタストアのプレフィックスを一覧表示します。 + +```bash +docker compose run --rm --entrypoint /bin/sh create-bucket -c \ + '/usr/bin/rc alias set rustfs http://rustfs:9000 "$RUSTFS_ACCESS_KEY" "$RUSTFS_SECRET_KEY" >/dev/null && /usr/bin/rc ls rustfs/my-bucket/trino-metastore --recursive' +``` + +```text +[2026-09-21 01:55:16] 155 B trino-metastore/.demo.trinoSchema +[2026-09-21 01:55:19] 474 B trino-metastore/demo/events/.trinoPermissions/user_trino +[2026-09-21 01:55:25] 1007 B trino-metastore/demo/events/.trinoSchema +[2026-09-21 01:55:25] 432 B trino-metastore/demo/events/20260921_..._cb761cec-...parquet +``` + +RustFS コンソールでこのプレフィックスを参照することもできます。 + +![RustFS コンソールに保存された Trino のメタデータとデータオブジェクト](./images/rustfs-trino-objects.png) + +## 6. RustFS S3 Tables を使用する + +RustFS S3 Tables は組み込みの Apache Iceberg REST カタログを提供します。Trino はテーブルバケットを管理された Iceberg ウェアハウスとして扱え、データは RustFS 内に保持されます。[S3 Tables](/administration/data/s3-tables) の説明に従ってテーブルバケットを有効化し、Trino の Iceberg コネクタを REST カタログに接続してください。REST カタログ URI は `http://:9000/iceberg`、ウェアハウスはバケット名で、カタログリクエスト(AWS Signature Version 4、署名名 `s3`)と S3 ファイルアクセスの両方がパススタイルのアドレス指定を使います。 + +S3 Tables のサポートマトリクスでは、Trino はカタログに対する読み取り専用プローブが行われた段階です。本番でこの経路を採用する前に、書き込み互換性と実際にデプロイする Trino バージョンを検証してください。 + +## 7. スタックを停止・リセットする + +RustFS データボリュームを保持したままコンテナを停止します。 + +```bash +docker compose down +``` + +保存したメタデータとデータを削除して空の RustFS ボリュームからやり直す場合は、明示的に `--volumes` を付けます。 + +```bash +docker compose down --volumes +``` + +## トラブルシューティング + +### `fs.native-s3.enabled` や `fs.s3.enabled` の設定エラー + +ネイティブ S3 ファイルシステムのプロパティ名は Trino のバージョン間で変わりました。Trino 435 は `fs.native-s3.enabled` を、それ以降のリリースは `fs.s3.enabled` を使用します。このガイドは `trinodb/trino:435` に固定しているため、`fs.native-s3.enabled` を使用してください。 + +### テーブル作成時に "Table directory must be ..." が出る + +ファイルメタストアでは、テーブルのロケーションは `hive.metastore.catalog.dir` の配下にある必要があります。スキーマの作成時に明示的なロケーションを指定せず、同じバケットプレフィックス内のディレクトリを指すようにしてください。 + +### Hive CSV ストレージ形式は VARCHAR のみ対応 + +CSV 形式は文字列以外の列を拒否します。型付きテーブルでは、このガイドのように `format = 'parquet'` を使用してください。 + +### AccessDenied や 403 レスポンス + +`hive.properties` の認証情報が RustFS の認証情報と一致しているか、`create-bucket` ジョブが正常に完了しているかを確認してください。 + +```bash +docker compose logs create-bucket +``` + +## 次のステップ + +- 追加の S3 オペレーションを採用する前に、[S3 互換性ノート](/administration/protocols/s3)を確認してください。 +- [アクセスキー管理](/security-compliance/iam/access-token)で本番用の専用認証情報を作成してください。 +- [Trino ドキュメント](https://trino.io/docs/current/)に従って、BI ツールを接続したりオブジェクトストレージカタログを追加したりしてください。 diff --git a/content/zh/developer/integration/big-data/flink.md b/content/zh/developer/integration/big-data/flink.md new file mode 100644 index 00000000..1547dcdf --- /dev/null +++ b/content/zh/developer/integration/big-data/flink.md @@ -0,0 +1,278 @@ +--- +title: "Apache Flink" +description: "使用 Apache Flink 及其 S3 文件系统插件读写 RustFS 对象存储中的 CSV 数据。" +--- + +本指南将 [Apache Flink](https://github.com/apache/flink) 通过 Flink 的 S3 文件系统插件(`flink-s3-fs-hadoop`)连接到 **RustFS**。你将使用 Docker Compose 启动一个 session 集群,以批模式把一个有界结果集写入桶中,再通过 Flink SQL 读回。整个流程使用 `flink:1.20` 和 `rustfs/rustfs-x86-musl:v2.3.1` 验证通过。 + +你需要安装带有 Compose 插件的 Docker。本部署用于本地集成测试,不适用于生产环境。 + +## 架构 + +```mermaid +flowchart LR + Job["Flink SQL job"] -->|"filesystem connector"| S3["S3 plugin (flink-s3-fs-hadoop)"] + S3 -->|"GET / PUT"| RustFS["RustFS :9000"] + Init["init-bucket job"] -->|"create my-bucket"| RustFS +``` + +`flink-s3-fs-hadoop` 插件为 Flink 的 filesystem 连接器注册了 `s3://` 协议。端点、path-style 寻址、纯 HTTP 和凭证通过 `flink-conf.yaml` 中的 `s3.*` 属性配置(经 `FLINK_PROPERTIES` 传入)。 + +## 1. 创建项目文件 + +创建工作目录: + +```bash +mkdir rustfs-flink +cd rustfs-flink +``` + +创建环境变量文件,并替换两个凭证占位符: + +```ini title=".env" +RUSTFS_ACCESS_KEY= +RUSTFS_SECRET_KEY= +``` + +请为桶使用专用的凭证,不要将 `.env` 提交到版本控制。 + +S3 插件随镜像内置在 `/opt/flink/opt/` 下,需要复制到 `/opt/flink/plugins/s3fs/` 才会加载。准备一个本地目录存放它: + +```bash +mkdir -p s3fs +docker create --name flink-tmp flink:1.20 +docker cp flink-tmp:/opt/flink/opt/flink-s3-fs-hadoop-1.20.5.jar s3fs/ +docker rm flink-tmp +``` + +创建 Compose 文件: + +```yaml title="compose.yaml" +services: + rustfs: + image: rustfs/rustfs-x86-musl:v2.3.1 + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + RUSTFS_VOLUMES: /data + RUSTFS_ADDRESS: ":9000" + RUSTFS_CONSOLE_ADDRESS: ":9001" + RUSTFS_CONSOLE_ENABLE: "true" + volumes: + - rustfs-data:/data + ports: + - "9000:9000" + - "9001:9001" + healthcheck: + test: ["CMD", "curl", "-sf", "http://127.0.0.1:9000/health"] + interval: 10s + timeout: 5s + retries: 6 + start_period: 10s + networks: + - flink + + create-bucket: + image: rustfs/rc:latest + depends_on: + rustfs: + condition: service_healthy + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + entrypoint: + - /bin/sh + - -c + - | + until /usr/bin/rc alias set rustfs http://rustfs:9000 "$${RUSTFS_ACCESS_KEY}" "$${RUSTFS_SECRET_KEY}"; do + echo "Waiting for RustFS..." + sleep 2 + done + /usr/bin/rc ls rustfs/my-bucket >/dev/null 2>&1 || /usr/bin/rc mb rustfs/my-bucket + networks: + - flink + + jobmanager: + image: flink:1.20 + command: jobmanager + environment: + FLINK_PROPERTIES: | + jobmanager.rpc.address: jobmanager + rest.address: jobmanager + rest.bind-address: 0.0.0.0 + s3.access-key: ${RUSTFS_ACCESS_KEY} + s3.secret-key: ${RUSTFS_SECRET_KEY} + s3.endpoint: http://rustfs:9000 + s3.path-style-access: true + volumes: + - ./s3fs:/opt/flink/plugins/s3fs:ro + ports: + - "8081:8081" + depends_on: + create-bucket: + condition: service_completed_successfully + networks: + - flink + + taskmanager: + image: flink:1.20 + command: taskmanager + environment: + FLINK_PROPERTIES: | + jobmanager.rpc.address: jobmanager + taskmanager.host: taskmanager + s3.access-key: ${RUSTFS_ACCESS_KEY} + s3.secret-key: ${RUSTFS_SECRET_KEY} + s3.endpoint: http://rustfs:9000 + s3.path-style-access: true + volumes: + - ./s3fs:/opt/flink/plugins/s3fs:ro + depends_on: + jobmanager: + condition: service_started + networks: + - flink + +networks: + flink: + +volumes: + rustfs-data: +``` + +`s3.access-key`、`s3.secret-key`、`s3.endpoint` 和 `s3.path-style-access` 属性为 JobManager 和 TaskManager 上的 S3 插件提供配置。 + +## 2. 启动部署 + +启动容器前先解析 Compose 文件: + +```bash +docker compose config +``` + +启动服务并等待桶初始化任务完成: + +```bash +docker compose up -d +docker compose ps -a +``` + +## 3. 把结果集写入 RustFS + +创建 SQL 作业——批模式加 filesystem sink: + +```yaml title="batch.sql" +SET 'execution.runtime-mode' = 'batch'; + +CREATE TABLE sink ( + id INT, + payload STRING +) WITH ( + 'connector' = 'filesystem', + 'path' = 's3://my-bucket/flink-out/', + 'format' = 'csv' +); + +INSERT INTO sink + VALUES (1, 'alpha'), (2, 'bravo'), (3, 'charlie'), (4, 'delta'), (5, 'echo'); +``` + +通过 JobManager 内的 SQL 客户端提交: + +```bash +docker compose exec jobmanager bash -c "/opt/flink/bin/sql-client.sh embedded -f /dev/stdin" < batch.sql +``` + +所有行写完后作业即结束。 + +## 4. 读回数据 + +创建读取查询——filesystem 连接器会扫描该前缀: + +```yaml title="read.sql" +CREATE TABLE readings ( + id INT, + payload STRING +) WITH ( + 'connector' = 'filesystem', + 'path' = 's3://my-bucket/flink-out/', + 'format' = 'csv' +); + +SET 'sql-client.execution.result-mode' = 'TABLEAU'; + +SELECT * FROM readings; +``` + +```bash +docker compose exec jobmanager bash -c "/opt/flink/bin/sql-client.sh embedded -f /dev/stdin" < read.sql +``` + +```text ++----+-------------+--------------------------------+ +| op | id | payload | ++----+-------------+--------------------------------+ +| +I | 1 | alpha | +| +I | 2 | bravo | +| +I | 3 | charlie | +| +I | 4 | delta | +| +I | 5 | echo | ++----+-------------+--------------------------------+ +``` + +## 5. 在 RustFS 中验证对象 + +通过桶初始化镜像列出该前缀: + +```bash +docker compose run --rm --entrypoint /bin/sh create-bucket -c \ + '/usr/bin/rc alias set rustfs http://rustfs:9000 "$RUSTFS_ACCESS_KEY" "$RUSTFS_SECRET_KEY" >/dev/null && /usr/bin/rc ls rustfs/my-bucket/flink-out --recursive' +``` + +```text +[2026-09-21 01:34:52] 41 B flink-out/part-f759f9e8-3d1b-46a1-a92e-53e9b727e831-task-0-file-0 +``` + +你也可以在 RustFS 控制台中浏览该前缀: + +![RustFS 控制台中存储的 Flink 输出文件](./images/rustfs-flink-objects.png) + +## 6. 使用 RustFS S3 Tables + +RustFS S3 Tables 提供内置的 Apache Iceberg REST 目录,Flink 可以把表桶当作托管的 Iceberg 仓库来使用,而数据仍保存在 RustFS 中。按 [S3 Tables](/administration/data/s3-tables) 的说明启用表桶,并把 Flink Iceberg 连接器的 REST 目录指向 RustFS:REST 目录 URI 为 `http://:9000/iceberg`,warehouse 即桶名,目录请求(AWS Signature Version 4,签名名 `s3`)与 S3 文件访问均使用 path-style 寻址。 + +根据 S3 Tables 支持矩阵,请在生产采用该路径前,用你实际部署的 Flink 和 Iceberg 版本完成验证。 + +## 7. 停止或重置环境 + +停止容器并保留 RustFS 数据卷: + +```bash +docker compose down +``` + +如需删除已存储的文件并从空的 RustFS 数据卷开始,请显式加上 `--volumes`: + +```bash +docker compose down --volumes +``` + +## 故障排除 + +### 写入时报 No AWS Credentials provided / AccessDenied + +S3 插件从 `flink-conf.yaml` 的 `s3.*` 属性读取凭证。确认 JobManager 和 TaskManager **两者**的 `FLINK_PROPERTIES` 中都包含 `s3.access-key`、`s3.secret-key`、`s3.endpoint` 和 `s3.path-style-access`,并且各自的 `/opt/flink/plugins/s3fs/` 中都存在插件 jar。 + +### TaskManager 无法解析 `rustfs` 主机名 + +所有 Flink 容器必须与 RustFS 共用一个 Compose 网络。如果 RustFS 挂在别的网络上,请把 Flink 容器也接入该网络后再提交作业。 + +### 流式写入在失败重试后报 "Stream closed" + +失败后恢复进行中的 S3 上传会让写入器进入无法恢复的状态。请删除桶中该作业的输出前缀后重新提交,或像本指南一样对一次性写入使用批模式。 + +## 后续步骤 + +- 在采用其他 S3 操作前,请查看 [S3 兼容性说明](/administration/protocols/s3)。 +- 通过[访问密钥管理](/security-compliance/iam/access-token)创建专用的生产凭证。 +- 阅读 [Apache Flink 文档](https://nightlies.apache.org/flink/flink-docs-stable/)了解 filesystem 连接器的分区与压实选项。 diff --git a/content/zh/developer/integration/big-data/images/rustfs-flink-objects.png b/content/zh/developer/integration/big-data/images/rustfs-flink-objects.png new file mode 100644 index 00000000..de2a82f3 Binary files /dev/null and b/content/zh/developer/integration/big-data/images/rustfs-flink-objects.png differ diff --git a/content/zh/developer/integration/big-data/images/rustfs-spark-objects.png b/content/zh/developer/integration/big-data/images/rustfs-spark-objects.png new file mode 100644 index 00000000..3bfc2aa0 Binary files /dev/null and b/content/zh/developer/integration/big-data/images/rustfs-spark-objects.png differ diff --git a/content/zh/developer/integration/big-data/images/rustfs-trino-objects.png b/content/zh/developer/integration/big-data/images/rustfs-trino-objects.png new file mode 100644 index 00000000..021aac1e Binary files /dev/null and b/content/zh/developer/integration/big-data/images/rustfs-trino-objects.png differ diff --git a/content/zh/developer/integration/big-data/index.md b/content/zh/developer/integration/big-data/index.md index fc1e56e7..c2b38258 100644 --- a/content/zh/developer/integration/big-data/index.md +++ b/content/zh/developer/integration/big-data/index.md @@ -12,5 +12,8 @@ description: "通过 S3 兼容的对象存储接口将数据分析系统连接 - [Milvus](./milvus.md) - [DuckDB](./duckdb.md) - [InfluxDB](./influxdb.md) +- [Spark](./spark.md) +- [Flink](./flink.md) +- [Trino](./trino.md) 将应用程序数据保存在专用存储桶和前缀中,并使用作用域限定为所需存储桶操作的凭证。 \ No newline at end of file diff --git a/content/zh/developer/integration/big-data/meta.json b/content/zh/developer/integration/big-data/meta.json index a03eefd4..448f23d8 100644 --- a/content/zh/developer/integration/big-data/meta.json +++ b/content/zh/developer/integration/big-data/meta.json @@ -5,6 +5,9 @@ "pyiceberg", "milvus", "duckdb", - "influxdb" + "influxdb", + "spark", + "flink", + "trino" ] } diff --git a/content/zh/developer/integration/big-data/spark.md b/content/zh/developer/integration/big-data/spark.md new file mode 100644 index 00000000..550b90b2 --- /dev/null +++ b/content/zh/developer/integration/big-data/spark.md @@ -0,0 +1,220 @@ +--- +title: "Apache Spark" +description: "通过 s3a 连接器读写 RustFS 对象存储中的 Parquet 数据。" +--- + +本指南将 [Apache Spark](https://github.com/apache/spark) 通过 `s3a` 连接器连接到 **RustFS**。你将使用 Docker Compose 启动 RustFS,运行一个把 Parquet 数据集写入桶中再读回的 Spark 作业,并在 RustFS 中验证这些对象。整个流程使用 `apache/spark:3.5.6`(配套 `hadoop-aws` 3.3.4)和 `rustfs/rustfs-x86-musl:v2.3.1` 验证通过。 + +你需要安装带有 Compose 插件的 Docker。本部署用于本地集成测试,不适用于生产环境。 + +## 架构 + +```mermaid +flowchart LR + Spark["Spark driver + executors"] -->|"S3AFileSystem"| RustFS["RustFS :9000"] + Init["init-bucket job"] -->|"create my-bucket"| RustFS +``` + +Spark 通过 `hadoop-aws` 的 S3A 文件系统访问 RustFS。端点、path-style 寻址、纯 HTTP 和凭证等连接器设置以 `spark.hadoop.fs.s3a.*` 属性的形式传入。 + +## 1. 创建项目文件 + +创建工作目录: + +```bash +mkdir rustfs-spark +cd rustfs-spark +``` + +创建环境变量文件,并替换两个凭证占位符: + +```ini title=".env" +RUSTFS_ACCESS_KEY= +RUSTFS_SECRET_KEY= +``` + +请为桶使用专用的凭证,不要将 `.env` 提交到版本控制。 + +创建 Spark 作业: + +```python title="job.py" +from pyspark.sql import SparkSession + +spark = SparkSession.builder.appName("rustfs-spark-demo").getOrCreate() +spark.sparkContext.setLogLevel("WARN") + +spark.range(1000).withColumnRenamed("id", "num") \ + .write.mode("overwrite").parquet("s3a://my-bucket/spark-demo/events") + +back = spark.read.parquet("s3a://my-bucket/spark-demo/events") +print("ROWS_READ_BACK:", back.count()) +spark.stop() +``` + +创建 Compose 文件: + +```yaml title="compose.yaml" +services: + rustfs: + image: rustfs/rustfs-x86-musl:v2.3.1 + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + RUSTFS_VOLUMES: /data + RUSTFS_ADDRESS: ":9000" + RUSTFS_CONSOLE_ADDRESS: ":9001" + RUSTFS_CONSOLE_ENABLE: "true" + volumes: + - rustfs-data:/data + ports: + - "9000:9000" + - "9001:9001" + healthcheck: + test: ["CMD", "curl", "-sf", "http://127.0.0.1:9000/health"] + interval: 10s + timeout: 5s + retries: 6 + start_period: 10s + networks: + - warehouse + + create-bucket: + image: rustfs/rc:latest + depends_on: + rustfs: + condition: service_healthy + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + entrypoint: + - /bin/sh + - -c + - | + until /usr/bin/rc alias set rustfs http://rustfs:9000 "$${RUSTFS_ACCESS_KEY}" "$${RUSTFS_SECRET_KEY}"; do + echo "Waiting for RustFS..." + sleep 2 + done + /usr/bin/rc ls rustfs/my-bucket >/dev/null 2>&1 || /usr/bin/rc mb rustfs/my-bucket + networks: + - warehouse + + spark: + image: apache/spark:3.5.6 + entrypoint: ["/opt/spark/bin/spark-submit"] + volumes: + - ./job.py:/job.py:ro + command: + - --conf + - spark.jars.ivy=/tmp/.ivy2 + - --packages + - org.apache.hadoop:hadoop-aws:3.3.4 + - --conf + - spark.hadoop.fs.s3a.endpoint=http://rustfs:9000 + - --conf + - spark.hadoop.fs.s3a.access.key=${RUSTFS_ACCESS_KEY} + - --conf + - spark.hadoop.fs.s3a.secret.key=${RUSTFS_SECRET_KEY} + - --conf + - spark.hadoop.fs.s3a.path.style.access=true + - --conf + - spark.hadoop.fs.s3a.connection.ssl.enabled=false + - --conf + - spark.hadoop.fs.s3a.impl=org.apache.hadoop.fs.s3a.S3AFileSystem + - /job.py + depends_on: + create-bucket: + condition: service_completed_successfully + networks: + - warehouse + +networks: + warehouse: + +volumes: + rustfs-data: +``` + +`--packages org.apache.hadoop:hadoop-aws:3.3.4` 会在启动时下载 S3 连接器,其版本必须与 Spark 镜像内置的 Hadoop 版本匹配。`spark.jars.ivy=/tmp/.ivy2` 把下载缓存移到可写目录。 + +## 2. 启动存储并运行作业 + +启动存储服务: + +```bash +docker compose up -d +docker compose ps -a +``` + +运行 Spark 作业: + +```bash +docker compose run --rm spark +``` + +```text +ROWS_READ_BACK: 1000 +``` + +## 3. 在 RustFS 中验证对象 + +通过桶初始化镜像列出数据集: + +```bash +docker compose run --rm --entrypoint /bin/sh create-bucket -c \ + '/usr/bin/rc alias set rustfs http://rustfs:9000 "$RUSTFS_ACCESS_KEY" "$RUSTFS_SECRET_KEY" >/dev/null && /usr/bin/rc ls rustfs/my-bucket/spark-demo --recursive' +``` + +```text +[2026-09-21 00:58:22] 0 B spark-demo/events/_SUCCESS +[2026-09-21 00:58:21] 1.46 KiB spark-demo/events/part-00000-...-c000.snappy.parquet +[2026-09-21 00:58:21] 1.46 KiB spark-demo/events/part-00001-...-c000.snappy.parquet +[2026-09-21 00:58:21] 1.46 KiB spark-demo/events/part-00002-...-c000.snappy.parquet +``` + +你也可以在 RustFS 控制台中浏览该前缀: + +![RustFS 控制台中存储的 Spark Parquet 输出](./images/rustfs-spark-objects.png) + +## 4. 使用 RustFS S3 Tables + +RustFS S3 Tables 提供内置的 Apache Iceberg REST 目录,Spark 可以把表桶当作托管的 Iceberg 仓库来使用,而数据仍保存在 RustFS 中。按 [S3 Tables](/administration/data/s3-tables) 的说明启用表桶并连接 Spark 的 Iceberg REST 目录:REST 目录 URI 为 `http://:9000/iceberg`,warehouse 即桶名,目录请求(AWS Signature Version 4,签名名 `s3`)与 S3 文件访问均使用 path-style 寻址。 + +根据 S3 Tables 支持矩阵,请在生产采用该路径前,用你实际部署的 Spark 和 Iceberg 版本完成验证。 + +## 5. 停止或重置环境 + +停止容器并保留 RustFS 数据卷: + +```bash +docker compose down +``` + +如需删除数据集并从空的 RustFS 数据卷开始,请显式加上 `--volumes`: + +```bash +docker compose down --volumes +``` + +## 故障排除 + +### NumberFormatException: For input string: "60s" + +`hadoop-aws` 版本与 Spark 镜像内置的 Hadoop 版本不匹配。Spark 4.x 镜像需要 `hadoop-aws` 3.4.x;本指南锁定 `apache/spark:3.5.6` 搭配 `hadoop-aws:3.3.4`。 + +### 连接 `rustfs:9000` 失败 + +`fs.s3a.endpoint` 是在 Compose 网络内解析的。如果 Spark 进程运行在宿主机上,请改用 (`http://localhost:9000`)。 + +### 返回 AccessDenied 或 403 响应 + +确认连接器设置与 RustFS 凭证一致,并确认 `create-bucket` 任务已成功完成: + +```bash +docker compose logs create-bucket +``` + +## 后续步骤 + +- 在采用其他 S3 操作前,请查看 [S3 兼容性说明](/administration/protocols/s3)。 +- 通过[访问密钥管理](/security-compliance/iam/access-token)创建专用的生产凭证。 +- 阅读 [Spark 文档](https://spark.apache.org/docs/latest/)了解结构化流与数据源选项。 diff --git a/content/zh/developer/integration/big-data/trino.md b/content/zh/developer/integration/big-data/trino.md new file mode 100644 index 00000000..efc51843 --- /dev/null +++ b/content/zh/developer/integration/big-data/trino.md @@ -0,0 +1,250 @@ +--- +title: "Trino" +description: "使用 Trino 与 hive 连接器的文件元存储,查询 RustFS 对象存储中的 CSV 和 Parquet 数据。" +--- + +本指南将 [Trino](https://github.com/trinodb/trino)——分布式 SQL 查询引擎——通过 hive 连接器的文件元存储和原生 S3 文件系统连接到 **RustFS**。你将创建 schema 和表、插入数据、查询回来,并在 RustFS 中验证这些对象。表的元数据和数据文件都保存在 RustFS 中。整个流程使用 `trinodb/trino:435` 和 `rustfs/rustfs-x86-musl:v2.3.1` 验证通过。 + +你需要安装带有 Compose 插件的 Docker。本部署用于本地集成测试,不适用于生产环境。 + +## 架构 + +```mermaid +flowchart LR + Client["trino CLI"] -->|"SQL"| Trino["Trino :8080"] + Trino -->|"metadata JSON"| RustFS["RustFS :9000"] + Trino -->|"data files"| RustFS + Init["init-bucket job"] -->|"create my-bucket"| RustFS +``` + +hive 连接器配合 `hive.metastore=file` 把 schema 和表元数据以 JSON 对象的形式保存在目录目录之下,原生 S3 文件系统(`fs.s3.enabled`)以纯 HTTP 上的 path-style 寻址把元数据和数据文件都存进 RustFS。 + +## 1. 创建项目文件 + +创建工作目录: + +```bash +mkdir rustfs-trino +cd rustfs-trino +``` + +创建环境变量文件,并替换两个凭证占位符: + +```ini title=".env" +RUSTFS_ACCESS_KEY= +RUSTFS_SECRET_KEY= +``` + +请为桶使用专用的凭证,不要将 `.env` 提交到版本控制。 + +创建 Trino 的 catalog 配置: + +```ini title="hive.properties" +connector.name=hive +hive.metastore=file +hive.metastore.catalog.dir=s3://my-bucket/trino-metastore +fs.s3.enabled=true +s3.endpoint=http://rustfs:9000 +s3.region=us-east-1 +s3.path-style-access=true +s3.aws-access-key= +s3.aws-secret-key= +``` + +`hive.metastore.catalog.dir` 把文件元存储指进桶内,元数据和数据都保存在 RustFS 中。`fs.s3.enabled` 启用原生 S3 文件系统;容器网络端点需要 `s3.path-style-access`。 + +创建 Compose 文件: + +```yaml title="compose.yaml" +services: + rustfs: + image: rustfs/rustfs-x86-musl:v2.3.1 + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + RUSTFS_VOLUMES: /data + RUSTFS_ADDRESS: ":9000" + RUSTFS_CONSOLE_ADDRESS: ":9001" + RUSTFS_CONSOLE_ENABLE: "true" + volumes: + - rustfs-data:/data + ports: + - "9000:9000" + - "9001:9001" + healthcheck: + test: ["CMD", "curl", "-sf", "http://127.0.0.1:9000/health"] + interval: 10s + timeout: 5s + retries: 6 + start_period: 10s + networks: + - warehouse + + create-bucket: + image: rustfs/rc:latest + depends_on: + rustfs: + condition: service_healthy + environment: + RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} + RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} + entrypoint: + - /bin/sh + - -c + - | + until /usr/bin/rc alias set rustfs http://rustfs:9000 "$${RUSTFS_ACCESS_KEY}" "$${RUSTFS_SECRET_KEY}"; do + echo "Waiting for RustFS..." + sleep 2 + done + /usr/bin/rc ls rustfs/my-bucket >/dev/null 2>&1 || /usr/bin/rc mb rustfs/my-bucket + networks: + - warehouse + + trino: + image: trinodb/trino:435 + volumes: + - ./hive.properties:/etc/trino/catalog/hive.properties:ro + - metastore-data:/data/metastore + depends_on: + create-bucket: + condition: service_completed_successfully + networks: + - warehouse + +networks: + warehouse: + +volumes: + rustfs-data: + metastore-data: +``` + +## 2. 启动部署 + +启动容器前先解析 Compose 文件: + +```bash +docker compose config +``` + +启动服务并等待桶初始化任务完成: + +```bash +docker compose up -d +docker compose ps -a +``` + +Trino 服务日志输出 `SERVER STARTED` 即表示已启动。容器以 `trino` 用户(uid 1000)运行,请确保元存储卷可写: + +```bash +docker compose exec trino id +docker compose exec trino ls -la /data/metastore +``` + +## 3. 创建 schema 与表 + +创建 schema 时不要指定 location——Trino 会把它放到 RustFS 中的目录目录之下: + +```bash +docker compose exec trino trino --execute \ + "CREATE SCHEMA hive.demo" +``` + +建表并插入五行数据: + +```bash +docker compose exec trino trino --execute \ + "CREATE TABLE hive.demo.events (id bigint, label varchar) WITH (format = 'parquet')" + +docker compose exec trino trino --execute \ + "INSERT INTO hive.demo.events VALUES (1,'alpha'),(2,'bravo'),(3,'charlie'),(4,'delta'),(5,'echo')" +``` + +```text +INSERT: 5 rows +``` + +## 4. 查询数据 + +把数据读回来: + +```bash +docker compose exec trino trino --execute \ + "SELECT * FROM hive.demo.events ORDER BY id" +``` + +```text +"1","alpha" +"2","bravo" +"3","charlie" +"4","delta" +"5","echo" +``` + +## 5. 在 RustFS 中验证对象 + +通过桶初始化镜像列出元存储前缀: + +```bash +docker compose run --rm --entrypoint /bin/sh create-bucket -c \ + '/usr/bin/rc alias set rustfs http://rustfs:9000 "$RUSTFS_ACCESS_KEY" "$RUSTFS_SECRET_KEY" >/dev/null && /usr/bin/rc ls rustfs/my-bucket/trino-metastore --recursive' +``` + +```text +[2026-09-21 01:55:16] 155 B trino-metastore/.demo.trinoSchema +[2026-09-21 01:55:19] 474 B trino-metastore/demo/events/.trinoPermissions/user_trino +[2026-09-21 01:55:25] 1007 B trino-metastore/demo/events/.trinoSchema +[2026-09-21 01:55:25] 432 B trino-metastore/demo/events/20260921_..._cb761cec-...parquet +``` + +你也可以在 RustFS 控制台中浏览该前缀: + +![RustFS 控制台中存储的 Trino 元数据与数据对象](./images/rustfs-trino-objects.png) + +## 6. 使用 RustFS S3 Tables + +RustFS S3 Tables 提供内置的 Apache Iceberg REST 目录,Trino 可以把表桶当作托管的 Iceberg 仓库来使用,而数据仍保存在 RustFS 中。按 [S3 Tables](/administration/data/s3-tables) 的说明启用表桶,并把 Trino 的 Iceberg 连接器接到 REST 目录:REST 目录 URI 为 `http://:9000/iceberg`,warehouse 即桶名,目录请求(AWS Signature Version 4,签名名 `s3`)与 S3 文件访问均使用 path-style 寻址。 + +根据 S3 Tables 支持矩阵,Trino 目前只对目录做过只读探测;在生产采用该路径前,请自行验证写入兼容性与实际部署的 Trino 版本。 + +## 7. 停止或重置环境 + +停止容器并保留 RustFS 数据卷: + +```bash +docker compose down +``` + +如需删除已存储的元数据和数据并从空的 RustFS 数据卷开始,请显式加上 `--volumes`: + +```bash +docker compose down --volumes +``` + +## 故障排除 + +### `fs.native-s3.enabled` 或 `fs.s3.enabled` 的配置错误 + +原生 S3 文件系统的属性名在不同 Trino 版本间有变化:Trino 435 使用 `fs.native-s3.enabled`,更新的版本使用 `fs.s3.enabled`。本指南锁定 `trinodb/trino:435`,请使用 `fs.native-s3.enabled`。 + +### 建表时报 "Table directory must be ..." + +文件元存储要求表位置位于 `hive.metastore.catalog.dir` 之下。创建 schema 时请不要指定 location,或把 schema location 指向同一桶前缀内的目录。 + +### Hive CSV 存储格式仅支持 VARCHAR + +CSV 格式不支持非字符串列。如本指南所示,有类型的表请使用 `format = 'parquet'`。 + +### 返回 AccessDenied 或 403 响应 + +确认 `hive.properties` 中的凭证与 RustFS 凭证一致,并确认 `create-bucket` 任务已成功完成: + +```bash +docker compose logs create-bucket +``` + +## 后续步骤 + +- 在采用其他 S3 操作前,请查看 [S3 兼容性说明](/administration/protocols/s3)。 +- 通过[访问密钥管理](/security-compliance/iam/access-token)创建专用的生产凭证。 +- 阅读 [Trino 文档](https://trino.io/docs/current/)连接 BI 工具并添加对象存储 catalog。