Showing posts with label AWS. Show all posts
Showing posts with label AWS. Show all posts

Mar 28, 2018

Integrating AWS SDK v2 with Akka streams

It's been a while since I last used AWS S3 a lot. In the mean time AWS came up with a new SDK and I started working with Scala and Akka full time. The SDK is still in beta but it's hard to imagine it'll take them more than a few more months to finish. Akka is still convenient but choke-full of APIs at different levels of complexity. Among other things I am researching how some of them could play together nicely. While digging, invariably a nugget or two can be found.

I wanted to plug S3 file download functionality from the SDK into a code path based on Akka streams. The former is all about CompletableFutures, the latter is about Source/Sink abstractions wired together with mostly built-in stages. Both are relatively large and self-sufficient frameworks with their own thread pools and coding idioms. I found one easy way to make them work together.

The trick is to notice that the AsyncResponseHandler returns a reactive stream Publisher. Yes, a third famous API in a file of less than a hundred lines. Akka can create a Source form a Publisher. It was also an unusual opportunity to use a Scala Promise (a quick reminder - in Scala a Promise provides write access to the state represented by a Future for read access). The idea is
  • as is custom in Scala, a time consuming operation is wrapped into a Future
  • a Promise is created before a download operation is started
  • when the download operation is started, the Promise is used to create a Future that is returned to the client
  • when the download operation finishes, the Promise is completed successfully
  • there is also the error case where the Promise is used to convey the error to the client

A unit test illustrates how a created Akka Source can be used to actually transfer some data using built-in Akka file Sink. I use a small S3 file made publicly available by some kind strangers.  

My first impression from this exercise is that it is reasonable but feels a little unnatural. Yes, the code itself is quite concise and the performance penalty must be negligible. But the j.u.c ecosystem is mostly alien to idiomatic Scala.  

May 22, 2016

Asynchronous AWS S3 file transfer

I am probably ten years too late in writing about the oldest AWS service but the power of AWS SDK asynchronous file transfer is irresistible. Moving a 300MB file takes just a few seconds. It is very easy to plug this API into reactive backend services. For this exercise we assume that one service uploads a file to S3 and another service periodically checks for new files in that location.

Fundamentally, asynchronous operations require an instance of the TransferManager class and a callback to process status notifications. I wrapped the whole process into a few classes representing abstractions for uploading to, downloading from, and detecting newly uploaded files in a pre-configured location in some S3 bucket location.

The TransferManager API typically takes a Request object and an asynchronous status listener. It returns a Transfer instance that can be used to retrieve error message in case of failure. Polling an S3 location for available files requires a loop because the results are returned in batches. Checking if a file exists at a given S3 path  is implemented as an attempt to fetch the corresponding file metadata and treating a thrown exception as "FileNotFound".

S3 is a simple (duh!) service so there are only two additional notes. First, it is a good idea to encode some metadata into file names on S3. Things such as tenant id or version or video resolution. It helps with deciding how to handle a downloaded file by parsing its name. Second, it's convenient to superimpose a "directory structure" onto the flat namespace of the S3 bucket abstraction. 

So a reasonable file naming convention might include a three-part prefix appended to all relative file paths: a backend service name, a file schema version, and a namespace representing either an environment (e.g. PROD) or a developer (in development deployments). The version part in particular makes upgrades much easier in production. For example,
"$BUCKET/prod/somesvc/v2/relative_path/file.ext".

The digram below shows a typical sequence of operations for uploading a file and then finding it with S3 polling from a different service.


In my example,
  • FileDownloader / FileUploader - abstractions used by the client to start a file transfer operation
  • TransferCallback - the callback interface called by file transfer operations to report final status to the client asynchronously
  • S3Destination - a way to specify a common bucket "subdirectory" for multiple files
  • S3TransferProgressListener - converts progress events to success/failure notifications; makes possible to extract an error message 
  • S3FileTransferClient / S3Client - TransferManager-based implementation of file transfer operations
  • S3Uploader / S3Downloader - S3Client-based implementation of file transfer API     
  • S3ChangeDetector - a job to be run periodically (e.g. with ScheduledExecutorService::scheduleWithFixedDelay) on the receiver side to look for new files on S3
  • FileHandler - the callback interface called by S3ChangeDetector for every found file not seen before (the way of keeping track of previously seen files is likely to be application-specific)


Mar 21, 2016

AWS SQS for reactive services

Last year I left the cozy world of self-managed systems for the now predominant AWS platform. It's been a mixed blessing so far. On the positive side, devops are relieved of so many daily troubles of running infrastructure components such as message brokers. On the negative side, developers are deprived of many conveniences provided by any JMS/AMQP-style product such as topics and no need for polling. 

It's very common to use a message broker to implement asynchronous API for your backend services. Especially in a world where major RPC frameworks (I am looking at you, PBs/Avro/Thrift) have only partial support for fully asynchronous clients and servers. Not everyone is running AKKA or Finagle in production. All you need is message broadcasting (i.e. topics) and an instance identifier to ignore responses sent to someone else.

By contrast, SQS is designed for essentially one topology: a number of worker nodes processing a shared queue. There is a gap between SQS and Kinesis which I still find to be rather painful. But for services with one-way requests only SQS is indeed simple to use.

For people coming from traditional message broker world there are a few noteworthy differences:
  • you need to poll a queue to recieve messages, there is no broker to push messages to you
  • you can receive up to ten messages in one request
  • there is a visibility timeout associated with messages; messages received by one consumer but not deleted from the queue within the visibility timeout are returned to the queue (and become available to other consumers)
  • an absolute queue name and a queue URL are two different strings
  • you can list queues by queue name prefix
I have a handy example illustrating SQS message processing with AWS SDK. It implements a typical message processing sequence as shown below:

  • SqsQueuePoller is supposed to be called periodically to poll a queue
  • AsyncSqsClient is a half-sync/half-async-style wrapper around AWS SDK client
  • Handler represents a service capable of processing multiple requests concurrently
  • MessageRepository keeps track of the messages being processed
  • VisibilityTimeoutTracker makes sure the visibility timeout never expires for the messages being processed