Add a BEAM read connector for JPA entities (#1132)

* Add a BEAM read connector for JPA entities

Added a Read connector to load JPA entities from Cloud SQL.

Also attempted a fix to the null threadfactory problem.
This commit is contained in:
Weimin Yu
2021-05-05 15:45:03 -04:00
committed by GitHub
parent 3f6ec8f1b0
commit b38574a9fc
3 changed files with 260 additions and 6 deletions
@@ -23,11 +23,18 @@ import com.google.common.collect.Streams;
import google.registry.backup.AppEngineEnvironment;
import google.registry.model.ofy.ObjectifyService;
import google.registry.persistence.transaction.JpaTransactionManager;
import google.registry.persistence.transaction.QueryComposer;
import google.registry.persistence.transaction.TransactionManagerFactory;
import java.io.Serializable;
import java.util.Objects;
import java.util.concurrent.ThreadLocalRandom;
import javax.persistence.criteria.CriteriaQuery;
import org.apache.beam.sdk.coders.Coder;
import org.apache.beam.sdk.coders.SerializableCoder;
import org.apache.beam.sdk.metrics.Counter;
import org.apache.beam.sdk.metrics.Metrics;
import org.apache.beam.sdk.transforms.Create;
import org.apache.beam.sdk.transforms.Deduplicate;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.GroupIntoBatches;
import org.apache.beam.sdk.transforms.PTransform;
@@ -36,6 +43,7 @@ import org.apache.beam.sdk.transforms.SerializableFunction;
import org.apache.beam.sdk.transforms.WithKeys;
import org.apache.beam.sdk.util.ShardedKey;
import org.apache.beam.sdk.values.KV;
import org.apache.beam.sdk.values.PBegin;
import org.apache.beam.sdk.values.PCollection;
/**
@@ -51,10 +59,143 @@ public final class RegistryJpaIO {
private RegistryJpaIO() {}
public static <R> Read<R, R> read(QueryComposerFactory<R> queryFactory) {
return Read.<R, R>builder().queryFactory(queryFactory).build();
}
public static <R, T> Read<R, T> read(
QueryComposerFactory<R> queryFactory, SerializableFunction<R, T> resultMapper) {
return Read.<R, T>builder().queryFactory(queryFactory).resultMapper(resultMapper).build();
}
public static <T> Write<T> write() {
return Write.<T>builder().build();
}
// TODO(mmuller): Consider detached JpaQueryComposer that works with any JpaTransactionManager
// instance, i.e., change composer.buildQuery() to composer.buildQuery(JpaTransactionManager).
// This way QueryComposer becomes reusable and serializable (at least with Hibernate), and this
// interface would no longer be necessary.
public interface QueryComposerFactory<T>
extends SerializableFunction<JpaTransactionManager, QueryComposer<T>> {}
/**
* A {@link PTransform transform} that executes a JPA {@link CriteriaQuery} and adds the results
* to the BEAM pipeline. Users have the option to transform the results before sending them to the
* next stages.
*
* <p>The BEAM pipeline may execute this transform multiple times due to transient failures,
* loading duplicate results into the pipeline. Before we add dedepuplication support, the easiest
* workaround is to map results to {@link KV} pairs, and apply the {@link Deduplicate} transform
* to the output of this transform:
*
* <pre>{@code
* PCollection<String> contactIds =
* pipeline
* .apply(RegistryJpaIO.read(
* (JpaTransactionManager tm) -> tm.createQueryComposer...,
* contact -> KV.of(contact.getRepoId(), contact.getContactId()))
* .withCoder(KvCoder.of(StringUtf8Coder.of(), StringUtf8Coder.of())))
* .apply(Deduplicate.keyedValues())
* .apply(Values.create());
* }</pre>
*/
@AutoValue
public abstract static class Read<R, T> extends PTransform<PBegin, PCollection<T>> {
public static final String DEFAULT_NAME = "RegistryJpaIO.Read";
abstract String name();
abstract RegistryJpaIO.QueryComposerFactory<R> queryFactory();
abstract SerializableFunction<R, T> resultMapper();
abstract TransactionMode transactionMode();
abstract Coder<T> coder();
abstract Builder<R, T> toBuilder();
@Override
public PCollection<T> expand(PBegin input) {
return input
.apply("Starting " + name(), Create.of((Void) null))
.apply(
"Run query for " + name(),
ParDo.of(new QueryRunner<>(queryFactory(), resultMapper())))
.setCoder(coder());
}
public Read<R, T> withName(String name) {
return toBuilder().name(name).build();
}
public Read<R, T> withResultMapper(SerializableFunction<R, T> mapper) {
return toBuilder().resultMapper(mapper).build();
}
public Read<R, T> withTransactionMode(TransactionMode transactionMode) {
return toBuilder().transactionMode(transactionMode).build();
}
public Read<R, T> withCoder(Coder<T> coder) {
return toBuilder().coder(coder).build();
}
static <R, T> Builder<R, T> builder() {
return new AutoValue_RegistryJpaIO_Read.Builder()
.name(DEFAULT_NAME)
.resultMapper(x -> x)
.transactionMode(TransactionMode.TRANSACTIONAL)
.coder(SerializableCoder.of(Serializable.class));
}
@AutoValue.Builder
public abstract static class Builder<R, T> {
abstract Builder<R, T> name(String name);
abstract Builder<R, T> queryFactory(RegistryJpaIO.QueryComposerFactory<R> queryFactory);
abstract Builder<R, T> resultMapper(SerializableFunction<R, T> mapper);
abstract Builder<R, T> transactionMode(TransactionMode transactionMode);
abstract Builder<R, T> coder(Coder coder);
abstract Read<R, T> build();
}
static class QueryRunner<R, T> extends DoFn<Void, T> {
private final QueryComposerFactory<R> querySupplier;
private final SerializableFunction<R, T> resultMapper;
QueryRunner(QueryComposerFactory<R> querySupplier, SerializableFunction<R, T> resultMapper) {
this.querySupplier = querySupplier;
this.resultMapper = resultMapper;
}
@ProcessElement
public void processElement(OutputReceiver<T> outputReceiver) {
// TODO(b/187210388): JpaTransactionManager should support non-transactional query.
// TODO(weiminyu): add deduplication
jpaTm()
.transactNoRetry(
() ->
querySupplier.apply(jpaTm()).stream()
.map(resultMapper::apply)
.forEach(outputReceiver::output));
// TODO(weiminyu): improve performance by reshuffle.
}
}
public enum TransactionMode {
NOT_TRANSACTIONAL,
TRANSACTIONAL;
}
}
/**
* A {@link PTransform transform} that writes a PCollection of entities to the SQL database using
* the {@link JpaTransactionManager}.