apache / apache/beam

Add declarative DSLs (XML & JSON)

Open
#17,969 0 comments 0 reactions 0 assignees View on GitHub
new feature P3 sdk-ideas
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
1d 20h
Merged PRs (30d)
196

Description

Even if users would still be able to use directly the API, it would be great to provide a DSL on top of the API covering batch and streaming data processing but also data integration.
Instead of designing a pipeline as a chain of apply() wrapping function (DoFn), we can provide a fluent DSL allowing users to directly leverage keyturn functions.

For instance, an user would be able to design a pipeline like:

```

.from(“kafka:localhost:9092?topic=foo”).reduce(...).split(...).wiretap(...).map(...).to(“jms:queue:foo….”);

```

The DSL will allow to use existing pipelines, for instance:

```

.from("cxf:...").reduce().pipeline("other").map().to("kafka:localhost:9092?topic=foo&acks=all")

```

So it means that we will have to create a IO Sink that can trigger the execution of a target pipeline: (from("trigger:other") triggering the pipeline execution when another pipeline design starts with pipeline("other")). We can also imagine to mix the runners: the pipeline() can be on one runner, the from("trigger:other") can be on another runner). It's not trivial, but it will give strong flexibility and key value for Beam.

In a second step, we can provide DSLs in different languages (the first one would be Java, but why not providing XML, akka, scala DSLs).

We can note in previous examples that the DSL would also provide data integration support to bean in addition of data processing. Data Integration is an extension of Beam API to support some Enterprise Integration Patterns (EIPs). As we would need metadata for data integration (even if metadata can also be interesting in stream/batch data processing pipeline), we can provide a DataxMessage built on top of PCollection. A DataxMessage would contain:
structured headers
binary payload
For instance, the headers can contains an Avro schema to describe the payload.
The headers can also contains useful information coming from the IO Source (for instance the partition/path where the data comes from, …).

Imported from Jira [BEAM-14](https://issues.apache.org/jira/browse/BEAM-14). Original Jira may contain additional context.
Reported by: jbonofre.

Contributor guide

Open the contributing guide

Research direction

No source files, tests, or concrete entry point are named. Start by reviewing Beam's existing pipeline API and the fluent examples in the issue, then determine the scope for batch, streaming, data integration, and cross-runner pipeline triggering. Done would require an agreed DSL design and implementation plan, not just a single edit.

Written by the indexing model from the issue text.

Assessment

Tech stack
java, json, xml
Domain
backend-api-design, data-engineering, stream-processing
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
18/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.