Skip to content

Commit 717b966

Browse files
committed
Use stable Kafka consumer group ID for job
1 parent 07b8bdf commit 717b966

2 files changed

Lines changed: 18 additions & 1 deletion

File tree

ingestion/src/main/java/feast/ingestion/ImportJob.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -120,6 +120,8 @@ public static PipelineResult runPipeline(ImportOptions options) throws IOExcepti
120120
"ReadFeatureRowFromSource",
121121
ReadFromSource.newBuilder()
122122
.setSource(source)
123+
.setConsumerGroupId(
124+
ReadFromSource.generateConsumerGroupId(source.getType(), store.getName()))
123125
.setSuccessTag(FEATURE_ROW_OUT)
124126
.setFailureTag(DEADLETTER_OUT)
125127
.build());

ingestion/src/main/java/feast/ingestion/transform/ReadFromSource.java

Lines changed: 16 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
import feast.ingestion.transform.fn.KafkaRecordToFeatureRowDoFn;
2424
import feast.storage.api.writer.FailedElement;
2525
import feast.types.FeatureRowProto.FeatureRow;
26+
import java.util.Optional;
2627
import org.apache.beam.sdk.io.kafka.KafkaIO;
2728
import org.apache.beam.sdk.transforms.PTransform;
2829
import org.apache.beam.sdk.transforms.ParDo;
@@ -41,6 +42,8 @@ public abstract class ReadFromSource extends PTransform<PBegin, PCollectionTuple
4142

4243
public abstract TupleTag<FailedElement> getFailureTag();
4344

45+
public abstract Optional<String> getConsumerGroupId();
46+
4447
public static Builder newBuilder() {
4548
return new AutoValue_ReadFromSource.Builder();
4649
}
@@ -54,6 +57,8 @@ public abstract static class Builder {
5457

5558
public abstract Builder setFailureTag(TupleTag<FailedElement> failureTag);
5659

60+
public abstract Builder setConsumerGroupId(String consumerGroupId);
61+
5762
abstract ReadFromSource autobuild();
5863

5964
public ReadFromSource build() {
@@ -83,7 +88,10 @@ public PCollectionTuple expand(PBegin input) {
8388
.withConsumerConfigUpdates(
8489
ImmutableMap.of(
8590
"group.id",
86-
generateConsumerGroupId(input.getPipeline().getOptions().getJobName())))
91+
getConsumerGroupId().isPresent()
92+
? getConsumerGroupId().get()
93+
: generateConsumerGroupId(
94+
input.getPipeline().getOptions().getJobName())))
8795
.withReadCommitted()
8896
.commitOffsetsInFinalize())
8997
.apply(
@@ -99,4 +107,11 @@ public PCollectionTuple expand(PBegin input) {
99107
private String generateConsumerGroupId(String jobName) {
100108
return "feast_import_job_" + jobName;
101109
}
110+
111+
/** Get a stable Kafka consumer group ID, so restarting jobs can resume offsets. */
112+
public static String generateConsumerGroupId(SourceType sourceType, String storeName) {
113+
return String.format("feast_import_job_%s_%s", sourceType, storeName)
114+
.replaceAll(" ", "_")
115+
.toLowerCase();
116+
}
102117
}

0 commit comments

Comments
 (0)