Avoiding broker roundtrip for recursive topology

Viewed 24

I have a tree of configs where each node can have arbitrary properties attached. The idea behind the tree is that properties of these tree nodes are derived from the parent node to the children nodes within the tree. Child properties with the same name overwrite the ones aggregated so far from the parent up to the root node. The goal is to generate aggregated maps of these config properties. Whenever any node in the tree gets an update all merged configs of its children and (grand children, and so on) shall be recalculated.

This is the topology I came up with and it works hoever I feel it's not optimal because the merged configs get first pushed up to the broker and then read back as input before being processed.

Is there a better way of achiving my goal? Is it possible to put all merged configs into a local state store and convert that into a KTable that I could use for the KTable-KTable-FK join?

Here is my code:

    data class Config(
        val id: String,
        val parentId: String?,
        var properties: Map<String, Any>,
        val eventDate: Int,
    )

    data class MergedConfig(
        val id: String,
        val parentId: String?,
        var properties: Map<String, Any>,
        val eventDate: Int,
    )


    private fun joinConfigWithParentConfig(): BiFunction<KStream<String, Config>, KTable<String, MergedConfig>, KStream<String, MergedConfig>> {
        return BiFunction<KStream<String, Config>, KTable<String, MergedConfig>, KStream<String, MergedConfig>> { configStream, mergedConfigTable ->
            val splitConfigStream =
                configStream
                    .repartition(Repartitioned.with(Serdes.String(), configSerde))
                    .split(Named.`as`("split-"))
                    .branch({ key, configEvent -> configEvent.parentId == null }, Branched.`as`("root"))
                    .branch({ key, configEvent -> configEvent.parentId != null }, Branched.`as`("children"))
                    .noDefaultBranch()

            // a stream, that will only contain the "merged" Config for the root
            val onlyRootStream = splitConfigStream["split-root"]!!
                .mapValues { key, value ->
                    MergedConfig(
                        value.id,
                        value.parentId,
                        value.properties,
                        time++
                    )
                }

            // a table that contains any incoming children within the tree, except the root
            val childrenTable = splitConfigStream["split-children"]!!.toTable()

            // joining only children with their merged parent config
            val mergedChildrenConfigTable = childrenTable.join(
                mergedConfigTable,
                { childConfig -> childConfig.parentId },
                { childConfig, mergedParentConfig ->
                    MergedConfig(
                        childConfig.id,
                        childConfig.parentId,
                        mergedParentConfig.properties + childConfig.properties, // merge config maps
                        time++
                    )
                }
            )

            // merge root and mergedChildren
            val mergedConfigStream = mergedChildrenConfigTable
                .toStream()
                .merge(onlyRootStream)
                .repartition(Repartitioned.with(Serdes.String(), mergedConfigSerde))

            mergedConfigStream
        }
    }
0 Answers
Related