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
RESOLVER_IDis the name the resolver is registered under;BUCKETandKEYname the object to read.fetchCredentials()builds aCredentialsobject withCredentials.builder(), settingACCESS_KEYandSECRET_KEYand an expiry one hour ahead withsetExpiresOn(). The expiry is what lets the cache know when a refresh is due.- A
SuppliedCredentialsResolverwraps the method reference and aCachingCredentialsResolverdecorates 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. CredentialsResolverRegistry.getSystemRegistry().add(RESOLVER_ID, resolver)registers the resolver in the JVM-wide registry.- The
AmazonS3FileSystemis configured withrequire(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. s3.open()resolves the credentials through the cache and connects;readFile()returns the object as anInputStreamfor aCSVReader.Job.run()transfers the records to aStreamWriteron the console ands3.close()in thefinallyblock disconnects.
Console Output
Each record in trades.csv is printed to the console, followed by the record count.
