Cache and Register a Credentials Resolver

This example combines two building blocks for handling credentials in production. A CachingCredentialsResolver caches the credentials returned by a slower delegate and calls it again only when the cached credentials are about to expire, when its time-to-live has elapsed or when refresh() is called. A CredentialsResolverRegistry maps an id to a resolver so that code, configuration and serialized pipelines carry a reference to the credentials instead of the credentials themselves.

The fetchCredentials() method stands in for a call to a secrets manager and returns keys that expire in one hour. Replace its body with the real lookup; the rest of the job stays the same.

Java Code Listing

package com.northconcepts.datapipeline.examples.security;

import java.io.InputStreamReader;
import java.util.concurrent.TimeUnit;

import com.northconcepts.datapipeline.amazons3.AmazonS3FileSystem;
import com.northconcepts.datapipeline.core.DataReader;
import com.northconcepts.datapipeline.core.DataWriter;
import com.northconcepts.datapipeline.core.StreamWriter;
import com.northconcepts.datapipeline.csv.CSVReader;
import com.northconcepts.datapipeline.job.Job;
import com.northconcepts.datapipeline.security.CachingCredentialsResolver;
import com.northconcepts.datapipeline.security.Credentials;
import com.northconcepts.datapipeline.security.CredentialsResolver;
import com.northconcepts.datapipeline.security.CredentialsResolverRegistry;
import com.northconcepts.datapipeline.security.SuppliedCredentialsResolver;

public class CacheAndRegisterACredentialsResolver {

    private static final String RESOLVER_ID = "s3-trades";
    private static final String BUCKET = "YOUR BUCKET";
    private static final String KEY = "output/trades.csv";

    public static void main(String[] args) throws Throwable {
        CredentialsResolver resolver = new CachingCredentialsResolver(
                new SuppliedCredentialsResolver(CacheAndRegisterACredentialsResolver::fetchCredentials),
                TimeUnit.MINUTES.toMillis(15));

        CredentialsResolverRegistry.getSystemRegistry().add(RESOLVER_ID, resolver);

        AmazonS3FileSystem s3 = new AmazonS3FileSystem()
                .setCredentialsResolver(CredentialsResolverRegistry.getSystemRegistry().require(RESOLVER_ID));
        s3.open();
        try {
            DataReader reader = new CSVReader(new InputStreamReader(s3.readFile(BUCKET, KEY)))
                    .setFieldNamesInFirstRow(true);
            DataWriter writer = StreamWriter.newSystemOutWriter();

            Job.run(reader, writer);
        } finally {
            s3.close();
        }
    }

    private static Credentials fetchCredentials() {
        return Credentials.builder()
                .set(Credentials.ACCESS_KEY, "YOUR ACCESS KEY")
                .set(Credentials.SECRET_KEY, "YOUR SECRET KEY")
                .setExpiresOn(System.currentTimeMillis() + TimeUnit.HOURS.toMillis(1))
                .build();
    }

}

Code Walkthrough

  1. RESOLVER_ID is the name the resolver is registered under; BUCKET and KEY name the object to read.
  2. fetchCredentials() builds a Credentials object with Credentials.builder(), setting ACCESS_KEY and SECRET_KEY and an expiry one hour ahead with setExpiresOn(). The expiry is what lets the cache know when a refresh is due.
  3. A SuppliedCredentialsResolver wraps the method reference and a CachingCredentialsResolver decorates it with a 15 minute time-to-live. The supplier is called on the first resolve and again only when the time-to-live passes or the cached credentials report that they are expiring.
  4. CredentialsResolverRegistry.getSystemRegistry().add(RESOLVER_ID, resolver) registers the resolver in the JVM-wide registry.
  5. The AmazonS3FileSystem is configured with require(RESOLVER_ID), which looks the resolver up by id and throws if nothing is registered under it. The S3 configuration only ever sees the id.
  6. s3.open() resolves the credentials through the cache and connects; readFile() returns the object as an InputStream for a CSVReader.
  7. Job.run() transfers the records to a StreamWriter on the console and s3.close() in the finally block disconnects.

Console Output

Each record in trades.csv is printed to the console, followed by the record count.

Mobile Analytics