Skip to content

Commit 9f841cc

Browse files
rajuGTrajuGT
andauthored
Dagger Custom Job Builder (#62)
With this commit, one can utilize the existing Dagger code, which supports existing source and sink adapters, serialization/deserialization, and monitoring setups. With just a minimal change (one config override), it can be integrated into the existing orchestrator/deployment cycle. We have also added an ExampleStreamApiJobBuilder as a reference implementation. It does not use Flink Table/SQL APIs; instead, it leverages the DataStream APIs to compute results. Co-authored-by: rajuGT <raju.gt@gojek.com>
1 parent 64d62da commit 9f841cc

8 files changed

Lines changed: 236 additions & 27 deletions

File tree

dagger-core/src/main/java/com/gotocompany/dagger/core/StreamManager.java renamed to dagger-core/src/main/java/com/gotocompany/dagger/core/DaggerSqlJobBuilder.java

Lines changed: 13 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -40,10 +40,7 @@
4040
import static com.gotocompany.dagger.functions.common.Constants.PYTHON_UDF_ENABLE_KEY;
4141
import static org.apache.flink.table.api.Expressions.$;
4242

43-
/**
44-
* The Stream manager.
45-
*/
46-
public class StreamManager {
43+
public class DaggerSqlJobBuilder implements JobBuilder {
4744

4845
private final Configuration configuration;
4946
private final StreamExecutionEnvironment executionEnvironment;
@@ -55,11 +52,11 @@ public class StreamManager {
5552
private final DaggerContext daggerContext;
5653

5754
/**
58-
* Instantiates a new Stream manager.
55+
* Instantiates dagger sql job-builder.
5956
*
60-
* @param daggerContext the daggerContext in form of param
57+
* @param daggerContext the daggerContext in form of param
6158
*/
62-
public StreamManager(DaggerContext daggerContext) {
59+
public DaggerSqlJobBuilder(DaggerContext daggerContext) {
6360
this.daggerContext = daggerContext;
6461
this.configuration = daggerContext.getConfiguration();
6562
this.executionEnvironment = daggerContext.getExecutionEnvironment();
@@ -71,7 +68,8 @@ public StreamManager(DaggerContext daggerContext) {
7168
*
7269
* @return the stream manager
7370
*/
74-
public StreamManager registerConfigs() {
71+
@Override
72+
public JobBuilder registerConfigs() {
7573
stencilClientOrchestrator = new StencilClientOrchestrator(configuration);
7674
org.apache.flink.configuration.Configuration flinkConfiguration = (org.apache.flink.configuration.Configuration) this.executionEnvironment.getConfiguration();
7775
daggerStatsDReporter = DaggerStatsDReporter.Provider.provide(flinkConfiguration, configuration);
@@ -96,7 +94,8 @@ public StreamManager registerConfigs() {
9694
*
9795
* @return the stream manager
9896
*/
99-
public StreamManager registerSourceWithPreProcessors() {
97+
@Override
98+
public JobBuilder registerSourceWithPreProcessors() {
10099
long watermarkDelay = configuration.getLong(Constants.FLINK_WATERMARK_DELAY_MS_KEY, Constants.FLINK_WATERMARK_DELAY_MS_DEFAULT);
101100
Boolean enablePerPartitionWatermark = configuration.getBoolean(Constants.FLINK_WATERMARK_PER_PARTITION_ENABLE_KEY, Constants.FLINK_WATERMARK_PER_PARTITION_ENABLE_DEFAULT);
102101
StreamsFactory.getStreams(configuration, stencilClientOrchestrator, daggerStatsDReporter)
@@ -144,7 +143,8 @@ private ApiExpression[] getApiExpressions(StreamInfo streamInfo) {
144143
*
145144
* @return the stream manager
146145
*/
147-
public StreamManager registerFunctions() throws IOException {
146+
@Override
147+
public JobBuilder registerFunctions() throws IOException {
148148
if (configuration.getBoolean(PYTHON_UDF_ENABLE_KEY, PYTHON_UDF_ENABLE_DEFAULT)) {
149149
PythonUdfConfig pythonUdfConfig = PythonUdfConfig.parse(configuration);
150150
PythonUdfManager pythonUdfManager = new PythonUdfManager(tableEnvironment, pythonUdfConfig, configuration);
@@ -178,7 +178,8 @@ private UdfFactory getUdfFactory(String udfFactoryClassName) throws ClassNotFoun
178178
*
179179
* @return the stream manager
180180
*/
181-
public StreamManager registerOutputStream() {
181+
@Override
182+
public JobBuilder registerOutputStream() {
182183
Table table = tableEnvironment.sqlQuery(configuration.getString(Constants.FLINK_SQL_QUERY_KEY, Constants.FLINK_SQL_QUERY_DEFAULT));
183184
StreamInfo streamInfo = createStreamInfo(table);
184185
streamInfo = addPostProcessor(streamInfo);
@@ -191,6 +192,7 @@ public StreamManager registerOutputStream() {
191192
*
192193
* @throws Exception the exception
193194
*/
195+
@Override
194196
public void execute() throws Exception {
195197
executionEnvironment.execute(configuration.getString(Constants.FLINK_JOB_ID_KEY, Constants.FLINK_JOB_ID_DEFAULT));
196198
}
Lines changed: 156 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,156 @@
1+
package com.gotocompany.dagger.core;
2+
3+
import com.gotocompany.dagger.common.configuration.Configuration;
4+
import com.gotocompany.dagger.common.core.DaggerContext;
5+
import com.gotocompany.dagger.common.core.StencilClientOrchestrator;
6+
import com.gotocompany.dagger.common.core.StreamInfo;
7+
import com.gotocompany.dagger.common.watermark.LastColumnWatermark;
8+
import com.gotocompany.dagger.common.watermark.StreamWatermarkAssigner;
9+
import com.gotocompany.dagger.common.watermark.WatermarkStrategyDefinition;
10+
import com.gotocompany.dagger.core.metrics.reporters.statsd.DaggerStatsDReporter;
11+
import com.gotocompany.dagger.core.processors.PreProcessorFactory;
12+
import com.gotocompany.dagger.core.processors.telemetry.processor.MetricsTelemetryExporter;
13+
import com.gotocompany.dagger.core.processors.types.Preprocessor;
14+
import com.gotocompany.dagger.core.sink.SinkOrchestrator;
15+
import com.gotocompany.dagger.core.source.StreamsFactory;
16+
import com.gotocompany.dagger.core.utils.Constants;
17+
import org.apache.flink.streaming.api.CheckpointingMode;
18+
import org.apache.flink.streaming.api.datastream.DataStream;
19+
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
20+
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
21+
import org.apache.flink.table.api.TableSchema;
22+
import org.apache.flink.types.Row;
23+
import org.apache.flink.util.Preconditions;
24+
25+
import java.io.IOException;
26+
import java.util.HashMap;
27+
import java.util.List;
28+
import java.util.Map;
29+
30+
public class ExampleStreamApiJobBuilder implements JobBuilder {
31+
32+
// static final String KEY_PATH = "meta.customer.id";
33+
34+
private final String inputStreamName1 = "data_streams_0";
35+
private final String inputStreamName2 = "data_streams_1";
36+
private final Map<String, StreamInfo> dataStreams = new HashMap<>();
37+
38+
private final DaggerContext daggerContext;
39+
private final Configuration configuration;
40+
private final StreamExecutionEnvironment executionEnvironment;
41+
private StencilClientOrchestrator stencilClientOrchestrator;
42+
private DaggerStatsDReporter daggerStatsDReporter;
43+
private final MetricsTelemetryExporter telemetryExporter = new MetricsTelemetryExporter();
44+
45+
public ExampleStreamApiJobBuilder(DaggerContext daggerContext) {
46+
this.daggerContext = daggerContext;
47+
this.configuration = daggerContext.getConfiguration();
48+
this.executionEnvironment = daggerContext.getExecutionEnvironment();
49+
}
50+
51+
@Override
52+
public JobBuilder registerConfigs() {
53+
stencilClientOrchestrator = new StencilClientOrchestrator(configuration);
54+
org.apache.flink.configuration.Configuration flinkConfiguration = (org.apache.flink.configuration.Configuration) this.executionEnvironment.getConfiguration();
55+
daggerStatsDReporter = DaggerStatsDReporter.Provider.provide(flinkConfiguration, configuration);
56+
57+
executionEnvironment.setMaxParallelism(configuration.getInteger(Constants.FLINK_PARALLELISM_MAX_KEY, Constants.FLINK_PARALLELISM_MAX_DEFAULT));
58+
executionEnvironment.getCheckpointConfig().setTolerableCheckpointFailureNumber(Integer.MAX_VALUE);
59+
executionEnvironment.enableCheckpointing(configuration.getLong(Constants.FLINK_CHECKPOINT_INTERVAL_MS_KEY, Constants.FLINK_CHECKPOINT_INTERVAL_MS_DEFAULT));
60+
executionEnvironment.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
61+
62+
// goes on...
63+
executionEnvironment.getConfig().setGlobalJobParameters(configuration.getParam());
64+
return this;
65+
}
66+
67+
@Override
68+
public JobBuilder registerSourceWithPreProcessors() {
69+
long watermarkDelay = configuration.getLong(Constants.FLINK_WATERMARK_DELAY_MS_KEY, Constants.FLINK_WATERMARK_DELAY_MS_DEFAULT);
70+
Boolean enablePerPartitionWatermark = configuration.getBoolean(Constants.FLINK_WATERMARK_PER_PARTITION_ENABLE_KEY, Constants.FLINK_WATERMARK_PER_PARTITION_ENABLE_DEFAULT);
71+
72+
StreamsFactory.getStreams(configuration, stencilClientOrchestrator, daggerStatsDReporter)
73+
.forEach(stream -> {
74+
String tableName = stream.getStreamName();
75+
76+
WatermarkStrategyDefinition watermarkStrategyDefinition = new LastColumnWatermark();
77+
78+
DataStream<Row> dataStream = stream.registerSource(executionEnvironment, watermarkStrategyDefinition.getWatermarkStrategy(watermarkDelay));
79+
StreamWatermarkAssigner streamWatermarkAssigner = new StreamWatermarkAssigner(new LastColumnWatermark());
80+
81+
DataStream<Row> dataStream1 = streamWatermarkAssigner
82+
.assignTimeStampAndWatermark(dataStream, watermarkDelay, enablePerPartitionWatermark);
83+
84+
85+
// just some legacy objects to adopt preprocessors
86+
TableSchema tableSchema = TableSchema.fromTypeInfo(dataStream.getType());
87+
StreamInfo streamInfo = new StreamInfo(dataStream1, tableSchema.getFieldNames());
88+
streamInfo = addPreProcessor(streamInfo, tableName);
89+
90+
if (tableName.equals(inputStreamName1)) {
91+
dataStreams.put(inputStreamName1, streamInfo);
92+
}
93+
if (tableName.equals(inputStreamName2)) {
94+
dataStreams.put(inputStreamName2, streamInfo);
95+
}
96+
});
97+
return this;
98+
}
99+
100+
@Override
101+
public JobBuilder registerFunctions() throws IOException {
102+
return this;
103+
}
104+
105+
@Override
106+
public JobBuilder registerOutputStream() {
107+
// NOTE - GET THE DATASTREAM REFERENCE
108+
StreamInfo streamInfo = dataStreams.get(inputStreamName1);
109+
Preconditions.checkNotNull(streamInfo, "Expected page log stream to be registered with name %s", inputStreamName1);
110+
111+
DataStream<Row> inputStream = streamInfo.getDataStream();
112+
113+
SinkOrchestrator sinkOrchestrator = new SinkOrchestrator(telemetryExporter);
114+
sinkOrchestrator.addSubscriber(telemetryExporter);
115+
116+
SingleOutputStreamOperator<Row> outputStream =
117+
inputStream
118+
119+
// NOTE - USE THE FLINK STREAM APIS HERE AND SINK THE OUTPUT
120+
121+
// .keyBy(
122+
// new KeySelector<Row, Integer>() {
123+
// private KeyExtractor keyExtractor;
124+
//
125+
// @Override
126+
// public Integer getKey(Row row) {
127+
// if (keyExtractor == null) {
128+
// keyExtractor = new KeyExtractor(row, KEY_PATH);
129+
// }
130+
// int userId = keyExtractor.extract(row);
131+
// return userId % DAU_PARALLELISM;
132+
// }
133+
// })
134+
// .process(new ShardedDistinctUserCounter())
135+
// .keyBy(r -> 0) // move all the output to one operator to calculate aggregation of all
136+
// .process(new UserCounterAggregator());
137+
.keyBy(r -> 0)
138+
.max("someField");
139+
140+
outputStream.sinkTo(sinkOrchestrator.getSink(configuration, new String[]{"uniq_users"}, stencilClientOrchestrator, daggerStatsDReporter));
141+
return this;
142+
}
143+
144+
@Override
145+
public void execute() throws Exception {
146+
executionEnvironment.execute(configuration.getString(Constants.FLINK_JOB_ID_KEY, Constants.FLINK_JOB_ID_DEFAULT));
147+
}
148+
149+
private StreamInfo addPreProcessor(StreamInfo streamInfo, String tableName) {
150+
List<Preprocessor> preProcessors = PreProcessorFactory.getPreProcessors(daggerContext, tableName, telemetryExporter);
151+
for (Preprocessor preprocessor : preProcessors) {
152+
streamInfo = preprocessor.process(streamInfo);
153+
}
154+
return streamInfo;
155+
}
156+
}
Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
1+
package com.gotocompany.dagger.core;
2+
3+
import java.io.IOException;
4+
5+
/**
6+
* An interface derived from the publicly exposed methods of {@code DaggerSqlJobBuilder}
7+
* previously referred as StreamManager.
8+
* <p>
9+
* The {@code KafkaProtoSQLProcessor}, which serves as the program entry point,
10+
* initializes an instance for the given {@code JOB_BUILDER_FQCN} value.
11+
* Ensure that the job builder class is bundled with the program during any
12+
* subsequent build stages. If it is not, the system falls back to the
13+
* {@code DEFAULT_JOB_BUILDER_FQCN} class, i.e., {@code com.gotocompany.dagger.core.DaggerSqlJobBuilder}.
14+
* <p>
15+
* Additionally, the job builder class is expected to provide a constructor
16+
* that accepts a single parameter of type {@code DaggerContext}
17+
*/
18+
public interface JobBuilder {
19+
20+
JobBuilder registerConfigs();
21+
22+
JobBuilder registerSourceWithPreProcessors();
23+
24+
JobBuilder registerFunctions() throws IOException;
25+
26+
JobBuilder registerOutputStream();
27+
28+
void execute() throws Exception;
29+
}

dagger-core/src/main/java/com/gotocompany/dagger/core/KafkaProtoSQLProcessor.java

Lines changed: 21 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,11 +3,15 @@
33
import com.gotocompany.dagger.common.configuration.Configuration;
44
import com.gotocompany.dagger.core.config.ConfigurationProvider;
55
import com.gotocompany.dagger.core.config.ConfigurationProviderFactory;
6+
import com.gotocompany.dagger.functions.common.Constants;
67
import org.apache.flink.client.program.ProgramInvocationException;
78
import com.gotocompany.dagger.common.core.DaggerContext;
89

10+
import java.lang.reflect.Constructor;
911
import java.util.TimeZone;
1012

13+
import static com.gotocompany.dagger.functions.common.Constants.JOB_BUILDER_FQCN_KEY;
14+
1115
/**
1216
* Main class to run Dagger.
1317
*/
@@ -25,8 +29,9 @@ public static void main(String[] args) throws ProgramInvocationException {
2529
Configuration configuration = provider.get();
2630
TimeZone.setDefault(TimeZone.getTimeZone("UTC"));
2731
DaggerContext daggerContext = DaggerContext.init(configuration);
28-
StreamManager streamManager = new StreamManager(daggerContext);
29-
streamManager
32+
33+
JobBuilder jobBuilder = getJobBuilderInstance(daggerContext);
34+
jobBuilder
3035
.registerConfigs()
3136
.registerSourceWithPreProcessors()
3237
.registerFunctions()
@@ -37,4 +42,18 @@ public static void main(String[] args) throws ProgramInvocationException {
3742
throw new ProgramInvocationException(e);
3843
}
3944
}
45+
46+
private static JobBuilder getJobBuilderInstance(DaggerContext daggerContext) {
47+
String className = daggerContext.getConfiguration().getString(JOB_BUILDER_FQCN_KEY, Constants.DEFAULT_JOB_BUILDER_FQCN);
48+
try {
49+
Class<?> builderClazz = Class.forName(className);
50+
Constructor<?> builderClazzConstructor = builderClazz.getConstructor(DaggerContext.class);
51+
return (JobBuilder) builderClazzConstructor.newInstance(daggerContext);
52+
} catch (Exception e) {
53+
Exception wrapperException = new Exception("Unable to instantiate job builder class: <" + className + "> \n"
54+
+ "Instantiating default job builder com.gotocompany.dagger.core.DaggerSqlJobBuilder", e);
55+
wrapperException.printStackTrace();
56+
return new DaggerSqlJobBuilder(daggerContext);
57+
}
58+
}
4059
}

dagger-core/src/test/java/com/gotocompany/dagger/core/StreamManagerTest.java renamed to dagger-core/src/test/java/com/gotocompany/dagger/core/DaggerSqlJobBuilderTest.java

Lines changed: 12 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -38,9 +38,9 @@
3838

3939
@PrepareForTest(TableSchema.class)
4040
@RunWith(PowerMockRunner.class)
41-
public class StreamManagerTest extends DaggerContextTestBase {
41+
public class DaggerSqlJobBuilderTest extends DaggerContextTestBase {
4242

43-
private StreamManager streamManager;
43+
private DaggerSqlJobBuilder daggerSqlJobBuilder;
4444

4545
private String jsonArray = "[\n"
4646
+ " {\n"
@@ -117,12 +117,12 @@ public void setup() {
117117
when(schema.getFieldNames()).thenReturn(new String[0]);
118118
PowerMockito.mockStatic(TableSchema.class);
119119
when(TableSchema.fromTypeInfo(typeInformation)).thenReturn(schema);
120-
streamManager = new StreamManager(daggerContext);
120+
daggerSqlJobBuilder = new DaggerSqlJobBuilder(daggerContext);
121121
}
122122

123123
@Test
124124
public void shouldRegisterRequiredConfigsOnExecutionEnvironment() {
125-
streamManager.registerConfigs();
125+
daggerSqlJobBuilder.registerConfigs();
126126

127127
verify(streamExecutionEnvironment, Mockito.times(1)).setParallelism(1);
128128
verify(streamExecutionEnvironment, Mockito.times(1)).enableCheckpointing(30000);
@@ -140,32 +140,32 @@ public void shouldRegisterSourceWithPreprocessorsWithWaterMarks() {
140140
when(source.assignTimestampsAndWatermarks(any(WatermarkStrategy.class))).thenReturn(singleOutputStream);
141141
when(singleOutputStream.getType()).thenReturn(typeInformation);
142142

143-
StreamManagerStub streamManagerStub = new StreamManagerStub(daggerContext, new StreamInfo(dataStream, new String[]{}));
144-
streamManagerStub.registerConfigs();
145-
streamManagerStub.registerSourceWithPreProcessors();
143+
DaggerSqlJobBuilderStub daggerSqlJobBuilderStub = new DaggerSqlJobBuilderStub(daggerContext, new StreamInfo(dataStream, new String[]{}));
144+
daggerSqlJobBuilderStub.registerConfigs();
145+
daggerSqlJobBuilderStub.registerSourceWithPreProcessors();
146146

147147
verify(streamTableEnvironment, Mockito.times(1)).fromDataStream(any(), new ApiExpression[]{});
148148
}
149149

150150
@Test
151151
public void shouldCreateOutputStream() {
152-
StreamManagerStub streamManagerStub = new StreamManagerStub(daggerContext, new StreamInfo(dataStream, new String[]{}));
153-
streamManagerStub.registerOutputStream();
152+
DaggerSqlJobBuilderStub daggerSqlJobBuilderStub = new DaggerSqlJobBuilderStub(daggerContext, new StreamInfo(dataStream, new String[]{}));
153+
daggerSqlJobBuilderStub.registerOutputStream();
154154
verify(streamTableEnvironment, Mockito.times(1)).sqlQuery("");
155155
}
156156

157157
@Test
158158
public void shouldExecuteJob() throws Exception {
159-
streamManager.execute();
159+
daggerSqlJobBuilder.execute();
160160

161161
verify(streamExecutionEnvironment, Mockito.times(1)).execute("SQL Flink job");
162162
}
163163

164-
final class StreamManagerStub extends StreamManager {
164+
final class DaggerSqlJobBuilderStub extends DaggerSqlJobBuilder {
165165

166166
private final StreamInfo streamInfo;
167167

168-
private StreamManagerStub(DaggerContext daggerContext, StreamInfo streamInfo) {
168+
private DaggerSqlJobBuilderStub(DaggerContext daggerContext, StreamInfo streamInfo) {
169169
super(daggerContext);
170170
this.streamInfo = streamInfo;
171171
}

dagger-functions/src/main/java/com/gotocompany/dagger/functions/common/Constants.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,4 +42,7 @@ public class Constants {
4242
public static final String COS_REGION = "COS_REGION";
4343
public static final String DEFAULT_COS_REGION = "ap-jakarta";
4444
public static final String ENABLE_TKE_OIDC_PROVIDER = "ENABLE_TKE_OIDC_PROVIDER";
45+
46+
public static final String JOB_BUILDER_FQCN_KEY = "JOB_BUILDER_FQCN";
47+
public static final String DEFAULT_JOB_BUILDER_FQCN = "com.gotocompany.dagger.core.DaggerSqlJobBuilder";
4548
}

docs/docs/concepts/architecture.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@ files as provided are consumed in a single stream.
2020

2121
_**Dagger Core**_
2222

23-
- The core part of the dagger(StreamManager) has the following responsibilities. It works sort of as a controller for other components in the dagger.
23+
- The core part of the dagger(DaggerSqlJobBuilder) has the following responsibilities. It works sort of as a controller for other components in the dagger.
2424
- Configuration management.
2525
- Table registration.
2626
- Configuring Deserialization and Serialization of data.

version.txt

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1 +1 @@
1-
0.11.9
1+
0.12.0

0 commit comments

Comments
 (0)