Export files to Cloudflare R2 using Scala and the AWS SDK
Upload a local file to Cloudflare R2 with the AWS SDK for Java v2, then use Cats Effect and FS2 for cancellable asynchronous uploads. The complete program below supports synchronous, asynchronous, streaming, and parallel modes.
Why choose Cloudflare R2?
R2 offers an S3-compatible API, so Scala applications can use the Java SDK. Compatibility is not identical to Amazon S3; use R2’s endpoint and supported request options.
Set up your Scala project
Use JDK 21 and sbt. Create
project/build.properties with this pinned launcher version:
sbt.version=1.10.11
Save the following as 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
Create src/main/scala/ for the program below. The Netty and Apache dependencies supply the
asynchronous and synchronous HTTP clients.
Configure the AWS SDK client
Create an R2 bucket and R2 S3 credentials.
Supply R2_ACCOUNT_ID, R2_ACCESS_KEY_ID, R2_SECRET_ACCESS_KEY, and R2_BUCKET through your
process environment. The program uses an existing bucket; it does not create one.
Use Region.of("auto"), path-style addressing, and disabled chunked encoding as shown in
Cloudflare’s Java guide.
WHEN_REQUIRED avoids optional SDK upload checksum trailers. Both clients sign requests and close
through Cats Effect Resource.
Basic synchronous upload
The sync mode below uses RequestBody.fromFile on the Cats Effect blocking pool. Its blocking
call waits for completion or the SDK timeout before releasing the client. Use either asynchronous
mode when you need cooperative cancellation.
Asynchronous upload with cats effect
Save this complete program as src/main/scala/R2Upload.scala. Unlike a synchronous RequestBody,
AsyncRequestBody.fromFile supplies the asynchronous client with a streaming body and its size.
Cats Effect’s
CompletableFuture bridge
propagates failures and requests cancellation when the fiber is canceled.
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)
}
}
}
With the four environment variables set and an existing local hello.txt, compile and run one mode:
sbt compile
sbt "runMain R2Upload async hello.txt uploads/hello.txt"
Replace async with sync or fs2 to exercise the other single-file paths. A successful PUT
replaces any existing object at the same key. Cancellation or a lost response does not prove the
server rejected the upload; check the destination before deciding whether to retry.
Handling responses and errors
The program reports success only after the SDK completes the upload. Service failures return a nonzero exit code and an HTTP status; local file, configuration, and network failures also return a nonzero exit code. A canceled upload is not converted into success. Keep detailed SDK diagnostics in private application logs when adapting this example.
Retry logic with exponential back-off
The SDK’s standard retry strategy classifies retryable failures and applies backoff. The example caps each request at three attempts, with a two-minute overall timeout. Avoid wrapping it in another unrestricted retry loop.
Keep each source file unchanged until the upload finishes, including retries. The SDK file body can reopen the file, and this pinned FS2 publisher reruns its file stream for each subscription. Replaying an identical file to the same key is suitable here; a changing file or a one-shot stream needs a different retry design.
Upload many files in parallel
The parallel mode uses parTraverseN(4) to limit active uploads. An upload failure cancels
sibling effects and propagates out of the client scope. Objects already uploaded remain in the
bucket; this is not a transaction.
sbt "runMain R2Upload parallel hello.txt uploads/hello.txt report.pdf uploads/report.pdf"
Stream large files using fs2
The fs2 mode reads 64 KiB chunks, converts them to the ByteBuffer elements required by AWS, and
delegates backpressure to FS2’s Reactive Streams publisher. It supplies the measured byte length on
both the request and body. toUnicastPublisher returns a Resource; its use scope stays open
until the upload completes, fails, or is canceled.
The source must be a finished, unchanged local file. Streaming limits buffering; it does not enable multipart upload or resume interrupted transfers. This example rejects files larger than 5 GiB before sending them. For larger objects, use R2 multipart uploads.
Security best practices
Keep credentials outside source control, restrict the R2 token to the intended bucket, and use HTTPS. Choose object keys deliberately: an upload can overwrite an existing object. If files come from users, finish your validation before making them available for download.
Common errors and their fixes
| Failure | What to check |
|---|---|
| HTTP 403 | R2 credentials and token permissions; disabled chunked encoding; correct account endpoint. |
| HTTP 404 | The bucket exists and its name matches R2_BUCKET. |
| Signature mismatch | System clock, endpoint, region auto, and SDK request options. |
| Timeout or connection failure | Network reachability and configured SDK timeouts. |
| Local failure | Arguments, environment variables, readable source files, and the single-PUT size limit. |
Wrap-up
The SDK file body is the simplest asynchronous path; FS2 is useful when an existing stream pipeline needs to supply the bytes. Both paths now share client cleanup, error propagation, and explicit request lengths.
For managed exports after file processing, Transloadit’s 🤖 /cloudflare/store Robot writes results to R2.
