Exporta archivos a Cloudflare R2 con Scala y el AWS SDK
Sube un archivo local a Cloudflare R2 con el AWS SDK para Java v2 y luego usa Cats Effect y FS2 para realizar subidas asíncronas que puedas cancelar. El programa completo que aparece a continuación admite los modos síncrono, asíncrono, de streaming y en paralelo.
¿Por qué elegir Cloudflare R2?
R2 ofrece una API compatible con S3, por lo que las aplicaciones de Scala pueden usar el SDK de Java. La compatibilidad no es idéntica a la de Amazon S3; usa el endpoint de R2 y las opciones de solicitud que admite.
Configura tu proyecto de Scala
Usa JDK 21 y sbt. Crea
project/build.properties con esta versión fija del lanzador:
sbt.version=1.10.11
Guarda lo siguiente como 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
Crea src/main/scala/ para el programa que aparece a continuación. Las dependencias de
Netty y Apache proporcionan los clientes HTTP asíncrono y síncrono.
Configura el cliente del AWS SDK
Crea un bucket de R2 y credenciales de S3 para R2. Proporciona
R2_ACCOUNT_ID, R2_ACCESS_KEY_ID, R2_SECRET_ACCESS_KEY y
R2_BUCKET mediante las variables de entorno del proceso. El programa usa un bucket
existente; no crea uno.
Usa Region.of("auto"), el direccionamiento basado en rutas y la codificación por bloques
desactivada, como se muestra en la guía de Java de Cloudflare.
WHEN_REQUIRED evita los campos finales opcionales de suma de comprobación que el SDK
añade a las subidas. Ambos clientes firman las solicitudes y se cierran mediante
Resource de Cats Effect.
Subida síncrona básica
El modo sync que aparece a continuación usa RequestBody.fromFile
en el grupo de hilos para operaciones bloqueantes de Cats Effect. Su llamada bloqueante espera a que
la operación termine o se alcance el tiempo de espera del SDK antes de liberar el cliente. Usa
cualquiera de los modos asíncronos cuando necesites cancelación cooperativa.
Subida asíncrona con Cats Effect
Guarda este programa completo como src/main/scala/R2Upload.scala. A diferencia de un
RequestBody síncrono, AsyncRequestBody.fromFile proporciona al cliente asíncrono
un cuerpo en streaming y su tamaño. El
puente de CompletableFuture de Cats Effect
propaga los errores y solicita la cancelación cuando se cancela la fibra.
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)
}
}
}
Con las cuatro variables de entorno configuradas y un archivo local hello.txt
existente, compila y ejecuta uno de los modos:
sbt compile
sbt "runMain R2Upload async hello.txt uploads/hello.txt"
Reemplaza async por sync o
fs2 para probar las otras opciones de subida de un solo archivo. Un PUT
exitoso reemplaza cualquier objeto existente con la misma clave. Una cancelación o la pérdida de una
respuesta no demuestra que el servidor haya rechazado la subida; comprueba el destino antes de decidir
si vuelves a intentarlo.
Gestión de respuestas y errores
El programa informa del éxito solo después de que el SDK complete la subida. Los errores del servicio devuelven un código de salida distinto de cero y un estado HTTP; los errores de archivos locales, configuración y red también devuelven un código de salida distinto de cero. Una subida cancelada no se convierte en un resultado exitoso. Guarda los diagnósticos detallados del SDK en registros privados de la aplicación cuando adaptes este ejemplo.
Lógica de reintentos con espera exponencial
La estrategia estándar de reintentos del SDK clasifica los errores que permiten reintentos y aplica intervalos de espera. El ejemplo limita cada solicitud a tres intentos, con un tiempo de espera total de dos minutos. Evita envolverlo en otro bucle de reintentos sin restricciones.
Mantén cada archivo de origen sin cambios hasta que termine la subida, incluidos los reintentos. El cuerpo de archivo del SDK puede volver a abrir el archivo, y el publicador de esta versión fija de FS2 vuelve a ejecutar su flujo de archivo para cada suscripción. En este caso, es adecuado volver a enviar un archivo idéntico a la misma clave; un archivo que cambia o un flujo de un solo uso necesita un diseño de reintentos diferente.
Sube muchos archivos en paralelo
El modo parallel usa parTraverseN(4) para limitar las subidas
activas. Un error de subida cancela los efectos que se ejecutan en paralelo y se propaga fuera del
ámbito del cliente. Los objetos ya subidos permanecen en el bucket; no se trata de una transacción.
sbt "runMain R2Upload parallel hello.txt uploads/hello.txt report.pdf uploads/report.pdf"
Transfiere archivos grandes en streaming con fs2
El modo fs2 lee bloques de 64 KiB, los convierte en los elementos
ByteBuffer que requiere AWS y delega el control de contrapresión al publicador
Reactive Streams de FS2. Proporciona la longitud medida en bytes tanto en la solicitud como en el
cuerpo. toUnicastPublisher devuelve un Resource; su ámbito de
use permanece abierto hasta que la subida se complete, falle o se cancele.
El origen debe ser un archivo local terminado y sin cambios. El streaming limita el almacenamiento en búfer; no habilita las subidas multiparte ni la reanudación de transferencias interrumpidas. Este ejemplo rechaza los archivos de más de 5 GiB antes de enviarlos. Para objetos más grandes, usa las subidas multiparte de R2.
Buenas prácticas de seguridad
Mantén las credenciales fuera del control de versiones, restringe el token de R2 al bucket previsto y usa HTTPS. Elige las claves de los objetos con cuidado: una subida puede sobrescribir un objeto existente. Si los archivos provienen de usuarios, termina de validarlos antes de ponerlos a disposición para su descarga.
Errores comunes y sus soluciones
| Error | Qué comprobar |
|---|---|
| HTTP 403 | Credenciales de R2 y permisos del token; codificación por bloques desactivada; endpoint correcto de la cuenta. |
| HTTP 404 | Que el bucket exista y su nombre coincida con R2_BUCKET. |
| La firma no coincide | Reloj del sistema, endpoint, región auto y opciones de solicitud del SDK. |
| Tiempo de espera agotado o fallo de conexión | Conectividad de red y tiempos de espera configurados en el SDK. |
| Error local | Argumentos, variables de entorno, archivos de origen legibles y límite de tamaño para un único PUT. |
Conclusión
El cuerpo de archivo del SDK es la opción asíncrona más sencilla; FS2 es útil cuando un pipeline de streaming existente necesita proporcionar los bytes. Ambas opciones ahora comparten la liberación de recursos del cliente, la propagación de errores y las longitudes explícitas de las solicitudes.
Para exportaciones gestionadas después del procesamiento de archivos, el Robot 🤖 /cloudflare/store de Transloadit guarda los resultados en R2.
