I have an apache beam pipeline that reads from pubsub, enriches data using Redis and finally writes to pubsub. I am trying to write tests to test the enrichment Dofn which is a stateful DoFn. Here the internal state is acting as a near cache to reduce the calls to Redis. For instantiating my Redis client I am using a factory declared in PipelineOptions such as
@Default.InstanceFactory(RedisClientFactory.class)
RedisClient getRedisClient();
void setRedisClient(RedisClient client);
In theory, the above client should be a singleton for each worker. In my unit tests, I am trying to mock some stuff inside that redis client. My tests look like this -
//setup pipeline
TestStream<MetricsInstance> inputStream =
TestStream.create(...).advanceWatermarkToInfinity();
PCollection<MetricsInstance> enrichedDataStream = pipeline.apply(inputStream)
.apply(ParDo.of(new ConvertToKeyValuePairDoFn<>()))
.apply(ParDo.of(new EnrichMetricsInstanceDoFn()));
CommonPipelineOptions options = PipelineOptionsFactory.as(CommonPipelineOptions.class);
RedisClient redisClient = options.getRedisClient();
JedisPool jedisPool = Mockito.mock(JedisPool.class);
jedis = Mockito.mock(Jedis.class);
Mockito.when(jedisPool.getResource()).thenReturn(jedis);
redisClient.setPool(jedisPool);
... some stubbing code and finally the pipeline run
PAssert.that(enrichedDataStream).containsInAnyOrder(expectedDataStream);
pipeline.run(options);
When I try to run this test I am getting an error like this
java.lang.IllegalArgumentException: Failed to serialize and deserialize property 'redisClient' with value 'xxx.xxx.RedisClientImpl@529cfee5'
To make the framework not attempt to serialize the client I can add @JsonIgnore on the getRedisClient() in my Options class. But that causes the Redis instance to be recreated at some point and all my mocking and stubbing is lost. I want to know whats the best way to test such scenarios.