Dateien mit Scala und AWS SDK nach Cloudflare R2 exportieren
Laden Sie eine lokale Datei mit dem AWS SDK for Java v2 nach Cloudflare R2 hoch. Nutzen Sie danach Cats Effect und FS2 für abbrechbare asynchrone Uploads. Das vollständige Programm unten unterstützt synchrone, asynchrone, Streaming- und parallele Modi.
Warum Cloudflare R2 wählen?
R2 bietet eine S3-kompatible API, sodass Scala-Anwendungen das Java SDK nutzen können. Die Kompatibilität ist nicht mit Amazon S3 identisch; verwenden Sie den R2-Endpunkt und die unterstützten Anfrageoptionen.
Scala-Projekt einrichten
Verwenden Sie JDK 21 und sbt. Erstellen Sie
project/build.properties mit dieser festgelegten Launcher-Version:
sbt.version=1.10.11
Speichern Sie Folgendes als build.sbt:
scalaVersion := "2.13.16"
libraryDependencies ++= Seq(
"software.amazon.awssdk" % "s3" % "2.54.19",
"software.amazon.awssdk" % "netty-nio-client" % "2.54.19",
"software.amazon.awssdk" % "apache-client" % "2.54.19",
"org.typelevel" %% "cats-effect" % "3.6.1",
"co.fs2" %% "fs2-io" % "3.12.0",
"co.fs2" %% "fs2-reactive-streams" % "3.12.0"
)
Compile / run / fork := true
Erstellen Sie src/main/scala/ für das folgende Programm. Die Abhängigkeiten Netty und
Apache stellen die asynchronen und synchronen HTTP-Clients bereit.
AWS-SDK-Client konfigurieren
Erstellen Sie einen R2-Bucket und R2-S3-Zugangsdaten. Stellen Sie
R2_ACCOUNT_ID, R2_ACCESS_KEY_ID, R2_SECRET_ACCESS_KEY und
R2_BUCKET über die Umgebungsvariablen Ihres Prozesses bereit. Das Programm
verwendet einen vorhandenen Bucket; es erstellt keinen.
Verwenden Sie Region.of("auto"), pfadbasierte Adressierung und deaktiviertes Chunked
Encoding, wie in der Java-Anleitung von Cloudflare gezeigt.
WHEN_REQUIRED verhindert optionale Prüfsummen-Trailer des SDK bei Uploads.
Beide Clients signieren Anfragen und werden mithilfe der Resource von Cats
Effect geschlossen.
Einfacher synchroner Upload
Der folgende Modus sync verwendet RequestBody.fromFile im
Blocking-Pool von Cats Effect. Sein blockierender Aufruf wartet auf den Abschluss oder das
SDK-Zeitlimit, bevor er den Client freigibt. Verwenden Sie einen der asynchronen Modi, wenn Sie
kooperativen Abbruch benötigen.
Asynchroner Upload mit Cats Effect
Speichern Sie dieses vollständige Programm als src/main/scala/R2Upload.scala. Anders als ein
synchroner RequestBody stellt AsyncRequestBody.fromFile dem asynchronen Client
einen Streaming-Body und dessen Größe bereit. Die
CompletableFuture-Anbindung von Cats Effect reicht Fehler weiter
und fordert einen Abbruch an, wenn die Fiber abgebrochen wird.
import cats.effect.{ExitCode, IO, IOApp, Resource}
import cats.effect.implicits._
import cats.syntax.all._
import fs2.interop.reactivestreams._
import fs2.io.file.{Files, Flags, Path}
import java.net.URI
import java.nio.ByteBuffer
import java.time.Duration
import java.util.Optional
import org.reactivestreams.Subscriber
import software.amazon.awssdk.auth.credentials.{AwsBasicCredentials, StaticCredentialsProvider}
import software.amazon.awssdk.core.async.AsyncRequestBody
import software.amazon.awssdk.core.checksums.RequestChecksumCalculation
import software.amazon.awssdk.core.client.config.ClientOverrideConfiguration
import software.amazon.awssdk.core.sync.RequestBody
import software.amazon.awssdk.regions.Region
import software.amazon.awssdk.retries.StandardRetryStrategy
import software.amazon.awssdk.services.s3.{S3AsyncClient, S3Client, S3Configuration}
import software.amazon.awssdk.services.s3.model.{PutObjectRequest, PutObjectResponse, S3Exception}
object R2Upload extends IOApp {
final case class Config(endpoint: URI, accessKey: String, secretKey: String, bucket: String)
final case class Task(path: Path, key: String)
private val service = S3Configuration.builder()
.pathStyleAccessEnabled(true)
.chunkedEncodingEnabled(false)
.build()
private def overrides: ClientOverrideConfiguration =
ClientOverrideConfiguration.builder()
.apiCallTimeout(Duration.ofMinutes(2))
.apiCallAttemptTimeout(Duration.ofSeconds(45))
.retryStrategy(StandardRetryStrategy.builder().maxAttempts(3).build())
.build()
private def credentials(config: Config): StaticCredentialsProvider =
StaticCredentialsProvider.create(
AwsBasicCredentials.create(config.accessKey, config.secretKey)
)
def asyncClient(config: Config): Resource[IO, S3AsyncClient] =
Resource.fromAutoCloseable(IO.blocking {
S3AsyncClient.builder()
.endpointOverride(config.endpoint)
.credentialsProvider(credentials(config))
.region(Region.of("auto"))
.serviceConfiguration(service)
.requestChecksumCalculation(RequestChecksumCalculation.WHEN_REQUIRED)
.overrideConfiguration(overrides)
.build()
})
def syncClient(config: Config): Resource[IO, S3Client] =
Resource.fromAutoCloseable(IO.blocking {
S3Client.builder()
.endpointOverride(config.endpoint)
.credentialsProvider(credentials(config))
.region(Region.of("auto"))
.serviceConfiguration(service)
.requestChecksumCalculation(RequestChecksumCalculation.WHEN_REQUIRED)
.overrideConfiguration(overrides)
.build()
})
private def request(bucket: String, task: Task, size: Long): PutObjectRequest = {
require(size >= 0 && size <= 5L * 1024 * 1024 * 1024,
"Use multipart upload for files larger than 5 GiB")
PutObjectRequest.builder()
.bucket(bucket).key(task.key)
.contentType("application/octet-stream")
.contentLength(size)
.build()
}
def uploadSync(client: S3Client, bucket: String, task: Task): IO[PutObjectResponse] =
Files[IO].size(task.path).flatMap { size =>
IO.blocking(client.putObject(
request(bucket, task, size), RequestBody.fromFile(task.path.toNioPath)
))
}
def uploadAsync(client: S3AsyncClient, bucket: String, task: Task): IO[PutObjectResponse] =
Files[IO].size(task.path).flatMap { size =>
IO.fromCompletableFuture(IO {
client.putObject(request(bucket, task, size),
AsyncRequestBody.fromFile(task.path.toNioPath))
})
}
def uploadStream(client: S3AsyncClient, bucket: String, task: Task): IO[PutObjectResponse] =
Files[IO].size(task.path).flatMap { size =>
val buffers = Files[IO].readAll(task.path, 64 * 1024, Flags.Read)
.chunks.map(chunk => ByteBuffer.wrap(chunk.toArray))
buffers.toUnicastPublisher.use { publisher =>
val body = new AsyncRequestBody {
override def contentLength(): Optional[java.lang.Long] =
Optional.of(java.lang.Long.valueOf(size))
override def subscribe(subscriber: Subscriber[_ >: ByteBuffer]): Unit =
publisher.subscribe(subscriber)
}
IO.fromCompletableFuture(IO {
client.putObject(request(bucket, task, size), body)
})
}
}
def uploadMany(client: S3AsyncClient, bucket: String,
tasks: List[Task]): IO[List[PutObjectResponse]] =
tasks.parTraverseN(4)(task => uploadAsync(client, bucket, task))
override def run(args: List[String]): IO[ExitCode] = {
val operation = IO {
require(args.length >= 3 && args.length % 2 == 1,
"Usage: R2Upload sync|async|fs2|parallel file key [file key ...]")
val mode = args.head
require(Set("sync", "async", "fs2", "parallel").contains(mode), "Unknown mode")
require(mode == "parallel" || args.length == 3, "Only parallel accepts multiple files")
val account = sys.env("R2_ACCOUNT_ID")
require(account.matches("[a-fA-F0-9]{32}"), "Invalid R2 account ID")
val config = Config(
URI.create(s"https://$account.r2.cloudflarestorage.com"),
sys.env("R2_ACCESS_KEY_ID"), sys.env("R2_SECRET_ACCESS_KEY"), sys.env("R2_BUCKET")
)
val tasks = args.tail.grouped(2).map(pair => Task(Path(pair.head), pair(1))).toList
(mode, config, tasks)
}.flatMap { case (mode, config, tasks) =>
val upload = if (mode == "sync")
syncClient(config).use(client => uploadSync(client, config.bucket, tasks.head).void)
else
asyncClient(config).use { client =>
mode match {
case "fs2" => uploadStream(client, config.bucket, tasks.head).void
case "parallel" => uploadMany(client, config.bucket, tasks).void
case _ => uploadAsync(client, config.bucket, tasks.head).void
}
}
upload *> IO.println("Upload completed.")
}
operation.as(ExitCode.Success).handleErrorWith {
case error: S3Exception =>
IO.println(s"Upload failed (HTTP ${error.statusCode()}).").as(ExitCode.Error)
case _ =>
IO.println("Upload failed. Check arguments, environment, file access, and network.")
.as(ExitCode.Error)
}
}
}
Wenn die vier Umgebungsvariablen gesetzt sind und die lokale Datei hello.txt
vorhanden ist, kompilieren Sie das Programm und führen Sie einen Modus aus:
sbt compile
sbt "runMain R2Upload async hello.txt uploads/hello.txt"
Ersetzen Sie async durch sync oder
fs2, um die anderen Varianten für einzelne Dateien auszuprobieren.
Ein erfolgreicher PUT ersetzt jedes vorhandene Objekt unter demselben Schlüssel. Ein Abbruch oder
eine verlorene Antwort beweist nicht, dass der Server den Upload abgelehnt hat. Prüfen Sie das Ziel,
bevor Sie entscheiden, ob Sie den Upload erneut versuchen.
Antworten und Fehler behandeln
Das Programm meldet Erfolg erst, nachdem das SDK den Upload abgeschlossen hat. Dienstfehler liefern einen Exit-Code ungleich null und einen HTTP-Status; Fehler bei lokalen Dateien, der Konfiguration und im Netzwerk liefern ebenfalls einen Exit-Code ungleich null. Ein abgebrochener Upload wird nicht als Erfolg gewertet. Bewahren Sie detaillierte SDK-Diagnosen in privaten Anwendungslogs auf, wenn Sie dieses Beispiel anpassen.
Wiederholungslogik mit exponentieller Wartezeit
Die Standardstrategie für Wiederholungen des SDK klassifiziert Fehler, bei denen ein erneuter Versuch möglich ist, und wendet eine Wartezeit an. Das Beispiel begrenzt jede Anfrage auf drei Versuche mit einem Gesamtzeitlimit von zwei Minuten. Vermeiden Sie eine weitere uneingeschränkte Wiederholungsschleife um diese Logik.
Lassen Sie jede Quelldatei bis zum Abschluss des Uploads unverändert, auch während erneuter Versuche. Der Datei-Body des SDK kann die Datei erneut öffnen, und dieser FS2-Publisher in der festgelegten Version führt seinen Dateistream für jedes Abonnement erneut aus. Eine identische Datei erneut unter demselben Schlüssel hochzuladen, ist hier geeignet. Eine sich ändernde Datei oder ein nur einmal lesbarer Stream benötigt eine andere Wiederholungsstrategie.
Viele Dateien parallel hochladen
Der Modus parallel begrenzt mit parTraverseN(4) die Zahl aktiver
Uploads. Ein Upload-Fehler bricht die anderen parallelen Effekte ab und wird aus dem
Gültigkeitsbereich des Clients weitergereicht. Bereits hochgeladene Objekte bleiben im Bucket;
dies ist keine Transaktion.
sbt "runMain R2Upload parallel hello.txt uploads/hello.txt report.pdf uploads/report.pdf"
Große Dateien mit fs2 streamen
Der Modus fs2 liest Blöcke von 64 KiB, wandelt sie in die von AWS benötigten
Elemente vom Typ ByteBuffer um und überlässt die Rückstaukontrolle dem
Reactive-Streams-Publisher von FS2. Er übergibt die gemessene Länge in Bytes sowohl in der Anfrage
als auch im Body. toUnicastPublisher gibt eine Resource zurück;
der zugehörige Gültigkeitsbereich von use bleibt offen, bis der Upload
abgeschlossen ist, fehlschlägt oder abgebrochen wird.
Die Quelle muss eine fertig geschriebene, unveränderte lokale Datei sein. Streaming begrenzt die Pufferung; es ermöglicht weder mehrteilige Uploads noch das Fortsetzen unterbrochener Übertragungen. Dieses Beispiel lehnt Dateien über 5 GiB vor dem Senden ab. Verwenden Sie für größere Objekte mehrteilige R2-Uploads.
Bewährte Verfahren für die Sicherheit
Bewahren Sie Zugangsdaten außerhalb der Versionsverwaltung auf, beschränken Sie das R2-Token auf den vorgesehenen Bucket und verwenden Sie HTTPS. Wählen Sie Objektschlüssel bewusst: Ein Upload kann ein vorhandenes Objekt überschreiben. Wenn Dateien von Nutzern stammen, schließen Sie Ihre Validierung ab, bevor Sie die Dateien zum Download bereitstellen.
Häufige Fehler und ihre Behebung
| Fehler | Zu prüfen |
|---|---|
| HTTP 403 | R2-Zugangsdaten und Token-Berechtigungen; deaktiviertes Chunked Encoding; korrekter Kontoendpunkt. |
| HTTP 404 | Der Bucket existiert und sein Name stimmt mit R2_BUCKET überein. |
| Signatur stimmt nicht überein | Systemuhr, Endpunkt, Region auto und SDK-Anfrageoptionen. |
| Zeitlimit oder Verbindungsfehler | Netzwerkerreichbarkeit und konfigurierte SDK-Zeitlimits. |
| Lokaler Fehler | Argumente, Umgebungsvariablen, lesbare Quelldateien und die Größenbegrenzung für einen einzelnen PUT. |
Fazit
Der Datei-Body des SDK ist die einfachste asynchrone Variante; FS2 ist nützlich, wenn eine vorhandene Stream-Pipeline die Bytes bereitstellen muss. Beide Varianten nutzen nun dieselbe Client-Freigabe, Fehlerweitergabe und explizite Anfragelängen.
Für verwaltete Exporte nach der Dateiverarbeitung schreibt der Robot 🤖 /cloudflare/store von Transloadit die Ergebnisse nach R2.
