Fortsetzbaren Go-Dateidownloader mit parallelen Blöcken bauen
Erstellen Sie einen ausführbaren Downloader, der vollständige Blöcke zwischen Programmläufen speichert. Er benötigt eine vertrauenswürdige SHA-256-Prüfsumme, einen Ursprungsserver mit starkem ETag und Unterstützung für Byte-Bereiche sowie ein privates lokales Ausgabeverzeichnis auf einem Unix-ähnlichen System.
Warum einen eigenen Dateidownloader entwickeln?
Vier Worker laden unabhängige Blöcke herunter. Vollständige Blöcke bleiben bei einer Unterbrechung erhalten, während teilweise heruntergeladene Blöcke niemals als vollständig gelten. Die Zieldatei wird erst dann atomar ersetzt, wenn die vollständige Datei der erwarteten Prüfsumme entspricht.
HTTP-Bereichsanfragen und Teilinhalte verstehen
Ein Range-Header ist eine Anfrage, keine Garantie. Jeder Block muss den Status 206, exakt den angeforderten Content-Range, dasselbe starke ETag und exakt die erwartete Anzahl an Bytes erhalten. If-Range verhindert, dass unbemerkt unterschiedliche Objektversionen kombiniert werden; eine Ausweichantwort mit Status 200 gilt hier als Fehler.
Semantik von HTTP-Bereichen und Validatoren (RFC 9110)
Einfachen Dateidownload in Go implementieren
Verwenden Sie Go 1.22 oder neuer. Erstellen Sie ein leeres Projektverzeichnis und initialisieren Sie das Modul. Das Programm verwendet ausschließlich die Standardbibliothek.
mkdir range-downloader
cd range-downloader
go mod init range-downloader
go mod edit -go=1.22
Unterstützung für fortsetzbare Downloads ergänzen
Unter dem Ausgabepfad mit dem Zusatz „.parts“ werden URL, ETag, erwartete Prüfsumme, Länge und Blockgröße gespeichert. Ein erneuter Programmlauf akzeptiert diesen Zustand nur, wenn alle Felder übereinstimmen. Er verwendet vollständige Blockdateien der erwarteten Länge erneut und prüft den zusammengesetzten Inhalt anhand der vertrauenswürdigen Prüfsumme. Bei einer Abweichung ist ein neuer Ausgabepfad erforderlich; verwenden Sie keinen beschädigten Zustand erneut.
Paralleles Herunterladen von Blöcken implementieren
Ein fester Worker-Pool verarbeitet Aufträge aus einem ungepufferten Kanal. Jede Antwort wird über einen Kopierpuffer von 32 KiB in eine temporäre Blockdatei gestreamt. Der erste Worker-Fehler bricht die übrigen Anfragen ab, und der Koordinator wartet vor der Rückkehr auf jeden Worker.
Fortschrittsbalken mit Echtzeitaktualisierungen ergänzen
Diese Version gibt beim abschließenden Zusammensetzen die Byte-Anzahl aus, statt eine Abhängigkeit für einen Fortschrittsbalken hinzuzufügen. Heruntergeladene Bytes werden erst als verifiziert angezeigt, wenn die abschließende Prüfsumme übereinstimmt.
Fehlerbehandlung und Wiederholungslogik
Fehler führen zu einem Prozessende mit einem Exitcode ungleich null. Führen Sie nach einem vorübergehenden Fehler denselben Befehl erneut aus, um vollständige Blöcke wiederzuverwenden. Es gibt keine automatische Wiederholungsschleife, die bei einem Autorisierungsfehler oder einer geänderten Repräsentation immer weiter neue Versuche starten könnte. Ctrl-C bricht Netzwerkanfragen ab und lässt vollständige Blöcke verfügbar.
Leistung durch Verbindungspooling optimieren
Ein gemeinsam genutzter HTTP-Client verwendet Verbindungen erneut, mit vier inaktiven Verbindungen pro Host und einem Zeitlimit von zwei Minuten pro Anfrage. Die Anzahl der Worker begrenzt die Parallelität im Netzwerk; der Festplattenbedarf umfasst die gespeicherten Blöcke und während des Zusammensetzens eine zweite vollständige Kopie.
Vollständiges Beispiel: einen CLI-Downloader erstellen
Speichern Sie dieses vollständige Programm als main.go. Beziehen Sie die SHA-256-Prüfsumme unabhängig vom Download von einem vertrauenswürdigen Herausgeber. Verwenden Sie ein übergeordnetes Verzeichnis, das Sie kontrollieren; kein anderer Prozess sollte die Ausgabe oder gespeicherte Blöcke verändern.
package main
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"os"
"os/signal"
"path/filepath"
"strings"
"sync"
"time"
)
const chunkSize int64 = 4 << 20
type identity struct {
URL, ETag, SHA256 string
Size, ChunkSize int64
}
func probe(ctx context.Context, client *http.Client, url, digest string) (identity, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodHead, url, nil)
if err != nil {
return identity{}, err
}
req.Header.Set("Accept-Encoding", "identity")
resp, err := client.Do(req)
if err != nil {
return identity{}, err
}
defer resp.Body.Close()
tag := resp.Header.Get("ETag")
if resp.StatusCode != 200 || resp.ContentLength < 0 || resp.ContentLength > 1<<40 ||
len(tag) < 2 || !strings.HasPrefix(tag, "\"") || !strings.HasSuffix(tag, "\"") ||
resp.Header.Get("Content-Encoding") != "" {
return identity{}, errors.New("need a known size (at most 1 TiB), strong ETag and unencoded HEAD 200")
}
return identity{url, tag, digest, resp.ContentLength, chunkSize}, nil
}
func fetchChunk(ctx context.Context, client *http.Client, id identity, dir string, start int64) error {
end := min(start+id.ChunkSize, id.Size) - 1
name := filepath.Join(dir, fmt.Sprintf("%d.part", start))
if info, err := os.Lstat(name); err == nil {
if info.Mode().IsRegular() && info.Size() == end-start+1 {
return nil
}
return errors.New("invalid saved chunk; use a new output path")
} else if !errors.Is(err, os.ErrNotExist) {
return err
}
req, err := http.NewRequestWithContext(ctx, http.MethodGet, id.URL, nil)
if err != nil {
return err
}
req.Header.Set("Accept-Encoding", "identity")
req.Header.Set("Range", fmt.Sprintf("bytes=%d-%d", start, end))
req.Header.Set("If-Range", id.ETag)
resp, err := client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
expected := fmt.Sprintf("bytes %d-%d/%d", start, end, id.Size)
if resp.StatusCode != http.StatusPartialContent || resp.Header.Get("Content-Range") != expected ||
resp.Header.Get("ETag") != id.ETag || resp.Header.Get("Content-Encoding") != "" ||
(resp.ContentLength != -1 && resp.ContentLength != end-start+1) {
return errors.New("server rejected range or changed representation")
}
tmp, err := os.CreateTemp(dir, ".chunk-")
if err != nil {
return err
}
defer os.Remove(tmp.Name())
defer tmp.Close()
n, err := io.CopyBuffer(tmp, io.LimitReader(resp.Body, end-start+2), make([]byte, 32<<10))
if err != nil {
return err
}
if n != end-start+1 {
return errors.New("incorrect chunk length")
}
if err := tmp.Sync(); err != nil {
return err
}
if err := tmp.Close(); err != nil {
return err
}
if err := ctx.Err(); err != nil {
return err
}
return os.Rename(tmp.Name(), name)
}
func download(ctx context.Context, client *http.Client, url, output, digest string, workers int) error {
sum, err := hex.DecodeString(digest)
if err != nil || len(sum) != sha256.Size || workers < 1 || workers > 16 {
return errors.New("provide a SHA-256 hex digest and 1–16 workers")
}
digest = strings.ToLower(digest)
id, err := probe(ctx, client, url, digest)
if err != nil {
return err
}
dir := output + ".parts"
fresh := false
if err := os.Mkdir(dir, 0700); err == nil {
fresh = true
} else if !errors.Is(err, os.ErrExist) {
return err
}
info, err := os.Lstat(dir)
if err != nil {
return err
}
if !info.IsDir() || info.Mode().Perm()&0077 != 0 {
return errors.New("parts directory must be private")
}
lock := filepath.Join(dir, ".lock")
if err := os.Mkdir(lock, 0700); err != nil {
return errors.New("parts directory locked; another download may be active")
}
defer os.Remove(lock)
manifest := filepath.Join(dir, "identity.json")
if fresh {
data, err := json.Marshal(id)
if err != nil {
return err
}
if err := os.WriteFile(manifest, data, 0600); err != nil {
return err
}
} else {
data, err := os.ReadFile(manifest)
if err != nil {
return err
}
var saved identity
if err := json.Unmarshal(data, &saved); err != nil {
return err
}
if saved != id {
return errors.New("saved download identity changed; use a new output path")
}
}
ctx, cancel := context.WithCancel(ctx)
defer cancel()
var wg sync.WaitGroup
var once sync.Once
var firstErr error
jobs := make(chan int64)
for i := 0; i < workers; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for start := range jobs {
if ctx.Err() != nil {
return
}
if err := fetchChunk(ctx, client, id, dir, start); err != nil {
once.Do(func() { firstErr = err; cancel() })
return
}
}
}()
}
send:
for start := int64(0); start < id.Size; start += id.ChunkSize {
select {
case jobs <- start:
case <-ctx.Done():
break send
}
}
close(jobs)
wg.Wait()
if firstErr != nil {
return firstErr
}
if err := ctx.Err(); err != nil {
return err
}
tmp, err := os.CreateTemp(filepath.Dir(output), ".download-")
if err != nil {
return err
}
defer os.Remove(tmp.Name())
defer tmp.Close()
hash := sha256.New()
for start := int64(0); start < id.Size; start += id.ChunkSize {
if err := ctx.Err(); err != nil {
return err
}
part, err := os.Open(filepath.Join(dir, fmt.Sprintf("%d.part", start)))
if err != nil {
return err
}
n, copyErr := io.Copy(io.MultiWriter(tmp, hash), io.LimitReader(part, id.ChunkSize+1))
closeErr := part.Close()
if copyErr != nil {
return copyErr
}
if closeErr != nil {
return closeErr
}
if n != min(id.ChunkSize, id.Size-start) {
return errors.New("saved chunk length changed")
}
fmt.Fprintf(os.Stderr, "Assembled: %d/%d bytes\n", start+n, id.Size)
}
if hex.EncodeToString(hash.Sum(nil)) != digest {
return errors.New("SHA-256 mismatch; use a new output path")
}
if err := tmp.Sync(); err != nil {
return err
}
if err := tmp.Close(); err != nil {
return err
}
if err := ctx.Err(); err != nil {
return err
}
return os.Rename(tmp.Name(), output)
}
func main() {
if len(os.Args) != 4 {
fmt.Fprintln(os.Stderr, "usage: downloader URL OUTPUT SHA256")
os.Exit(2)
}
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt)
defer stop()
transport := http.DefaultTransport.(*http.Transport).Clone()
transport.MaxIdleConnsPerHost = 4
defer transport.CloseIdleConnections()
client := &http.Client{Transport: transport, Timeout: 2 * time.Minute}
if err := download(ctx, client, os.Args[1], os.Args[2], os.Args[3], 4); err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
}
Kompilieren Sie das Programm und übergeben Sie anschließend Ihre Download-URL, den Ausgabepfad und die erwartete Prüfsumme über Umgebungsvariablen. Diese Prüfungen schlagen sofort fehl, wenn ein erforderlicher Wert fehlt.
GOTOOLCHAIN=local go build -o downloader .
./downloader "${DOWNLOAD_URL:?Set DOWNLOAD_URL}" "${OUTPUT_PATH:?Set OUTPUT_PATH}" "${EXPECTED_SHA256:?Set EXPECTED_SHA256}"
Bewährte Verfahren und häufige Fallstricke
Diese Implementierung lehnt unbekannte Längen, schwache oder fehlende ETags, codierte Antworten und Objekte größer als 1 TiB ab. Leere Dateien werden unterstützt, wenn HEAD ein starkes ETag und eine Länge von null liefert. SHA-256 erkennt beschädigte gespeicherte Blöcke und vermischte Versionen, selbst wenn sich ein Ursprungsserver fehlerhaft verhält.
Nach erfolgreichen Programmläufen bleibt das private Verzeichnis für die Blöcke erhalten, damit der Betreiber es ausdrücklich bereinigen kann. Ein erzwungenes Prozessende kann darin das Verzeichnis .lock zurücklassen: Entfernen Sie diese Sperre erst, nachdem Sie sich vergewissert haben, dass kein Downloader aktiv ist. Für das atomare Umbenennen müssen die temporäre zusammengesetzte Datei und das Ziel im selben Dateisystem liegen; das Ersetzungsverhalten ist hier auf Unix-ähnliche Systeme ausgelegt. Dies garantiert keine dauerhafte Speicherung bei einem Stromausfall.
Fazit
Sie haben nun einen vollständigen Downloader für Byte-Bereiche mit begrenzter Parallelität, dauerhaft gespeichertem Fortsetzungszustand und Verifizierung vor der Veröffentlichung.
Der Robot /http/import von Transloadit übernimmt HTTP-Importe in Verarbeitungsabläufen.
