Nebenläufige Dateisynchronisation mit Go und MinIO
Das Synchronisieren von Dateien aus Ihrem Object Storage in Ihre lokale Umgebung erfordert häufig, viele Dateien effizient zu verarbeiten. In diesem DevTip zeigen wir, wie Sie mit MinIO, einem quelloffenen, S3-kompatiblen Cloud-Storage-Server, ein einfaches, nebenläufiges Werkzeug zur Dateisynchronisation in Go bauen.
MinIO-Server einrichten
Bevor wir mit dem Code beginnen, richten wir mit Docker einen lokalen MinIO-Server ein. So erhalten wir eine Entwicklungsumgebung, in der wir unsere Implementierung testen können:
docker run \
-p 127.0.0.1:9000:9000 \
-p 127.0.0.1:9001:9001 \
--name minio \
-v ~/minio/data:/data \
-e "MINIO_ROOT_USER=minioadmin" \
-e "MINIO_ROOT_PASSWORD=minioadmin" \
quay.io/minio/minio server /data --console-address ":9001"
Dies ist eine Wegwerf-Sandbox mit den allgemein bekannten Zugangsdaten minioadmin/minioadmin, daher sind beide
Ports explizit an 127.0.0.1 gebunden. Ein bloßes -p 9000:9000 würde sie auf jeder
Schnittstelle des Hosts veröffentlichen, was auf einem Laptop meist auch das lokale Netzwerk
bedeutet.
Öffnen Sie die Web-Konsole unter http://localhost:9001, melden Sie sich mit diesen Zugangsdaten
an, erstellen Sie einen Bucket namens my-sync-bucket und laden Sie eine Handvoll Dateien
hinein. Der Code unten verwendet durchgehend diesen einen Bucket-Namen.
Projekt initialisieren
Erstellen Sie ein neues Go-Projekt und installieren Sie das MinIO-SDK. Die eigene Datei go.mod des gepinnten SDK gibt go 1.22 an, daher schlägt der Befehl go get unten mit einer älteren Toolchain fehl:
go mod init minio-sync
go get github.com/minio/minio-go/v7@v7.0.87
MinIO-Client einrichten
Erstellen Sie eine Verbindung zu Ihrer MinIO-Instanz mit sauberer Fehlerbehandlung. Der
Docker-Container oben spricht einfaches HTTP, daher steht useSSL hier auf false und die TLS-Einstellungen des Transports bleiben inaktiv.
Erst wenn Sie useSSL gegenüber einem Endpunkt mit TLS-Terminierung auf true umstellen, werden sie aktiv:
// client.go
package main
import (
"context"
"crypto/tls"
"fmt"
"net/http"
"time"
"github.com/minio/minio-go/v7"
"github.com/minio/minio-go/v7/pkg/credentials"
)
// bucketName is shared by every file in this program. Create it in the console first.
const bucketName = "my-sync-bucket"
func createMinioClient(ctx context.Context) (*minio.Client, error) {
endpoint := "localhost:9000"
accessKeyID := "minioadmin" // Use environment variables in production
secretAccessKey := "minioadmin" // Use environment variables in production
useSSL := false // true once your endpoint terminates TLS
// Only exercised when useSSL is true
transport := &http.Transport{
TLSClientConfig: &tls.Config{MinVersion: tls.VersionTLS12},
IdleConnTimeout: 90 * time.Second,
}
// Initialize minio client
opts := &minio.Options{
Creds: credentials.NewStaticV4(accessKeyID, secretAccessKey, ""),
Secure: useSSL,
Transport: transport,
}
client, err := minio.New(endpoint, opts)
if err != nil {
return nil, err
}
// BucketExists reports a bucket that is simply absent as (false, nil), so the boolean
// has to be checked too. Ignoring it turns a typo in the bucket name into a sync that
// quietly downloads nothing.
exists, err := client.BucketExists(ctx, bucketName)
if err != nil {
return nil, fmt.Errorf("failed to reach bucket %q: %w", bucketName, err)
}
if !exists {
return nil, fmt.Errorf("bucket %q does not exist", bucketName)
}
return client, nil
}
Nebenläufige Datei-Downloads implementieren
Hier ist eine verbesserte Implementierung mit sauberer Fehlerbehandlung, Kontextverwaltung und Aufräumlogik:
// sync.go
package main
import (
"context"
"fmt"
"io"
"os"
"path"
"path/filepath"
"strings"
"sync"
"time"
"github.com/minio/minio-go/v7"
)
type DownloadResult struct {
ObjectName string
Error error
}
func downloadFiles(ctx context.Context, client *minio.Client, bucketName string, outputDir string) error {
// 0700, because these are private copies. The process umask would otherwise usually
// leave them group- and world-readable.
if err := os.MkdirAll(outputDir, 0o700); err != nil {
return fmt.Errorf("failed to create output directory: %w", err)
}
// Resolve the directory once so the per-object checks below compare against a path
// that has no symlinks left in it.
root, err := filepath.EvalSymlinks(outputDir)
if err != nil {
return fmt.Errorf("failed to resolve output directory: %w", err)
}
// Cancelling here unblocks every worker if we bail out early
ctx, cancel := context.WithCancel(ctx)
defer cancel()
jobs := make(chan string, 100)
results := make(chan DownloadResult, 100)
var wg sync.WaitGroup
workerCount := 5
// Start workers
for i := 0; i < workerCount; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for objectName := range jobs {
err := downloadObject(ctx, client, bucketName, objectName, root)
select {
case results <- DownloadResult{ObjectName: objectName, Error: err}:
case <-ctx.Done():
return
}
}
}()
}
go func() {
wg.Wait()
close(results)
}()
// Collect results while the producer is still queueing. Draining only after
// close(jobs) deadlocks as soon as the buffers fill: workers block on
// `results <-`, stop reading `jobs`, and the producer blocks on `jobs <-`.
//
// Only a count and the first error are kept. A bucket can hold millions of objects,
// and appending every failure to a slice would grow without bound in exactly the run
// where things are going worst.
type failures struct {
count int
first error
}
collected := make(chan failures, 1)
go func() {
var seen failures
for result := range results {
if result.Error == nil {
continue
}
seen.count++
if seen.first == nil {
seen.first = fmt.Errorf("%s: %w", result.ObjectName, result.Error)
}
}
collected <- seen
}()
// List and queue objects. The deferred close lets the workers drain and exit
// on every path, including the early returns below.
listErr := func() error {
defer close(jobs)
opts := minio.ListObjectsOptions{Recursive: true}
for obj := range client.ListObjects(ctx, bucketName, opts) {
if obj.Err != nil {
// Cancel immediately so in-flight downloads stop now rather than running
// to completion behind a listing that has already failed.
cancel()
return fmt.Errorf("error listing objects: %w", obj.Err)
}
// Empty directory markers contain no file bytes to download.
if obj.Size == 0 && strings.HasSuffix(obj.Key, "/") {
continue
}
select {
case jobs <- obj.Key:
case <-ctx.Done():
return ctx.Err()
}
}
return nil
}()
// Always drain first, so every worker has finished before we report anything.
seen := <-collected
if listErr != nil {
return listErr
}
// A cancelled context ends the listing loop without an error of its own, and workers
// that were cut off never deliver a result. Without this check a cancelled sync would
// be indistinguishable from an empty bucket that synced cleanly.
if err := ctx.Err(); err != nil {
return fmt.Errorf("sync did not complete: %w", err)
}
if seen.count > 0 {
return fmt.Errorf("encountered %d download errors, first: %w", seen.count, seen.first)
}
return nil
}
// resolveOutputPath maps a server-supplied object key onto a path inside outputDir, which
// must already be symlink-free (see filepath.EvalSymlinks above).
//
// Keys are rejected rather than cleaned. `a/./b.txt` and `a/b/../b.txt` are different keys
// that both clean to `a/b.txt`, so cleaning them would let one object silently overwrite
// another, and lexical checks alone cannot see a symlink that is already on disk.
func resolveOutputPath(outputDir, objectName string) (string, error) {
if objectName == "" || objectName != path.Clean(objectName) || strings.Contains(objectName, `\`) {
return "", fmt.Errorf("refusing non-canonical object key: %q", objectName)
}
if !filepath.IsLocal(filepath.FromSlash(objectName)) {
return "", fmt.Errorf("refusing unsafe object key: %q", objectName)
}
current := outputDir
for _, element := range strings.Split(objectName, "/") {
current = filepath.Join(current, element)
info, err := os.Lstat(current)
if err != nil {
if os.IsNotExist(err) {
continue // Nothing here yet, so nothing can redirect the write
}
return "", fmt.Errorf("failed to inspect %q: %w", current, err)
}
if info.Mode()&os.ModeSymlink != 0 {
return "", fmt.Errorf("refusing object key through a symlink: %q", objectName)
}
}
return current, nil
}
func downloadObject(ctx context.Context, client *minio.Client, bucket, objectName, outputDir string) error {
outputPath, err := resolveOutputPath(outputDir, objectName)
if err != nil {
return err
}
// Create context with timeout
ctx, cancel := context.WithTimeout(ctx, 10*time.Minute)
defer cancel()
// Get object
obj, err := client.GetObject(ctx, bucket, objectName, minio.GetObjectOptions{})
if err != nil {
return fmt.Errorf("failed to get object: %w", err)
}
defer obj.Close()
if err := os.MkdirAll(filepath.Dir(outputPath), 0o700); err != nil {
return fmt.Errorf("failed to create directories: %w", err)
}
// Download to a sibling temp file so a failed sync never truncates the copy
// that a previous run completed, and never leaves a half-written file behind.
// os.CreateTemp already creates the file 0600.
temp, err := os.CreateTemp(filepath.Dir(outputPath), filepath.Base(outputPath)+".part-*")
if err != nil {
return fmt.Errorf("failed to create temporary file: %w", err)
}
tempPath := temp.Name()
defer func() {
temp.Close()
os.Remove(tempPath) // No-op once the rename below succeeded
}()
if _, err := io.Copy(temp, obj); err != nil {
return fmt.Errorf("failed to download file: %w", err)
}
if err := temp.Sync(); err != nil {
return fmt.Errorf("failed to flush file: %w", err)
}
if err := temp.Close(); err != nil {
return fmt.Errorf("failed to close file: %w", err)
}
// Rename is atomic within a filesystem, so readers see either the old file or
// the complete new one
if err := os.Rename(tempPath, outputPath); err != nil {
return fmt.Errorf("failed to publish file: %w", err)
}
return nil
}
Fehlerbehandlung und Fehlersuche
Leere Verzeichnis-Marker werden übersprungen. Objektschlüssel müssen auf sichere, kanonische
relative Pfade abgebildet werden, die das lokale Betriebssystem unterstützt; sie werden nicht
stillschweigend umbenannt. Zeitstempelnamen mit Doppelpunkten funktionieren beispielsweise unter
Unix, werden unter Windows aber von filepath.IsLocal abgelehnt.
Hier sind häufige Probleme, die auftreten können, und wie Sie sie beheben:
-
Verbindungsfehler:
- Prüfen Sie, ob der MinIO-Endpunkt erreichbar ist
- Überprüfen Sie die Firewall-Einstellungen
- Stellen Sie korrekte Zugangsdaten sicher
-
Berechtigungsprobleme:
- Prüfen Sie die Zugriffsrechte für den Bucket
- Überprüfen Sie die Dateisystemberechtigungen für das Ausgabeverzeichnis
-
Ressourcenbeschränkungen:
- Passen Sie die Anzahl der Worker an die Systemressourcen an
- Überwachen Sie den Speicherverbrauch bei großen Dateien
- Erwägen Sie, eine Ratenbegrenzung zu implementieren
Ihre Implementierung testen
Die drei Snippets oben bilden das gesamte Programm: client.go, sync.go und main.go unten,
nebeneinander im Modul minio-sync, das Sie zuvor erstellt haben. Führen Sie es mit go run . aus:
// main.go
package main
import (
"context"
"log"
"os/signal"
"syscall"
)
func main() {
// Ctrl-C cancels the context, which stops the workers and makes downloadFiles
// report the interruption instead of exiting as if the sync had finished.
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
// Create MinIO client
client, err := createMinioClient(ctx)
if err != nil {
log.Fatalf("Failed to create MinIO client: %v", err)
}
// Start download. bucketName comes from client.go.
if err := downloadFiles(ctx, client, bucketName, "./downloads"); err != nil {
log.Fatalf("Download failed: %v", err)
}
log.Println("Download completed successfully")
}
Sicherheitsüberlegungen
Für den Einsatz in der Produktion:
- Verwenden Sie Umgebungsvariablen oder einen sicheren Konfigurationsmanager für Zugangsdaten
- Aktivieren Sie TLS in Produktionsumgebungen
- Implementieren Sie geeignete Zugriffskontrollen für Buckets
- Verwenden Sie nach Möglichkeit temporäre Zugangsdaten
- Rotieren Sie Zugriffsschlüssel regelmäßig
Objektschlüssel stammen vom Server, daher behandelt resolveOutputPath sie als nicht
vertrauenswürdige Eingabe und lehnt alles ab, was absolut oder nicht kanonisch ist oder über einen
Symlink führt, der bereits im Ausgabeverzeichnis existiert. Das ist eine sinnvolle Prüfung, aber
keine Dateisystem-Sandbox: Der Code prüft den Pfad und öffnet ihn anschließend; dieses
Go-1.22-kompatible Beispiel schließt das Zeitfenster zwischen diesen beiden Schritten nicht. Richten
Sie outputDir auf ein Verzeichnis, das ausschließlich Ihrem Prozess gehört und
restriktive Berechtigungen hat, und nicht auf einen Baum, den andere lokale Benutzer oder Prozesse
während einer laufenden Synchronisation verändern können.
Abschließende Gedanken
Diese Implementierung bietet einen Ausgangspunkt für nebenläufige Datei-Downloads aus MinIO: Ergebnisse werden abgearbeitet, während die Auflistung noch läuft, unsichere Objektschlüssel werden abgelehnt, bevor sie das Dateisystem berühren, und jedes Objekt wird atomar veröffentlicht. Denken Sie daran, die Anzahl der Worker und die Timeout-Werte an Ihren konkreten Anwendungsfall und Ihre Systemressourcen anzupassen.
Einiges tut sie bewusst nicht: Sie prüft eine bereits synchronisierte Datei nie erneut (jeder Lauf lädt alles neu herunter), sie hat keine Ratenbegrenzung, das Timeout von 10 Minuten pro Objekt ist ein fester Wert und nicht von der Objektgröße abgeleitet, und Fehler werden als Anzahl plus erster Fehler zusammengefasst statt als vollständige Liste.
Viel Spaß beim Programmieren!
