Big Data 21 min read

Data Nirvana: AI Agents Orchestrate Trustworthy Data Pipelines

The article presents a unified data development platform where Forge Agent acts as an AI facade, Sindri handles deterministic planning via DSL, ADF manages execution and scheduling, and quality systems ensure trustworthiness, moving beyond chat interfaces to controlled, automated data production.

JD Tech Talk
JD Tech Talk
JD Tech Talk
Data Nirvana: AI Agents Orchestrate Trustworthy Data Pipelines

Introduction: AI in Production ≠ AI Taking Over Production

Current large models can understand requirements, generate code, and explain errors, but data production involves more than writing SQL. A requirement must pass through data source confirmation, definition clarification, processing orchestration, resource planning, task submission, scheduling, result delivery, and fault handling. Missing any step can turn a seemingly correct task into a production incident.

Therefore, the goal is not to let models bypass the platform to operate production directly, but to establish a controlled chain: Forge Agent as the unified AI facade, Sindri Plugin/MCP as deterministic tool development and constraints, Sindri backend for core chain planning and submission (including DSL parsing, execution plan generation, and optimization), and ADF to assist upstream systems in shielding big data platform operations while carrying DAG, scheduling, state, and lifecycle.

Why Build This System

Current business scenarios lack a unified, complete data development platform. Data ingestion, cleansing, business processing, index building, and operation maintenance remain highly manual, with fragmented and non-reusable development chains.

Main site : No dedicated platform carries the full data R&D process; ingestion, processing, and index building still rely on manual maintenance.

Recommendation : Has a relatively complete index building system, but high customization means ingestion, cleansing, and business processing must be done manually before the access system.

Vertical sites : Primarily use configuration control systems, but internal processes and processing steps differ across verticals, still requiring manual development for specific businesses.

These problems cause duplicate construction, high labor costs, and make it difficult to unify R&D standards, process handoffs, quality assurance, and operation maintenance. The solution is a unified data development system that connects ingestion, cleansing, business processing, index building, task submission, scheduling, and fault handling, preserving flexibility for different business scenarios while precipitating reusable platform capabilities to reduce duplicate development and manual operations, and improving data quality through systematic capabilities.

System Architecture Overview

The aim is not to add a chat entry to the old chain, but to let intelligent entry, domain semantics, unified base, and quality assurance each take their place and collaborate as a system.

Horizontal main chain : Forge Agent organizes the entire R&D process; the domain data R&D platform precipitates business semantics through two main lines — Feature Engineering (feature definition, computation selection, publishing) and Base Pool Processing (domain DSL, semantic compilation and planning, task submission and management). ADF serves as the unified data and compute base, carrying operator systems, connections and adapters, data and metadata, and multi-engine execution. The output ultimately serves search, recommendation, key scenarios, and online applications. Natural language intent flows through domain semantics, execution DAG, data products, and back to business feedback.

Vertical quality assurance : The data quality system does not belong to a single layer but runs vertically through the entire chain — data reconciliation, data quality, lineage tracking, and monitoring dashboards provide trust from four angles: consistency, compliance, traceability, and observability.

Forge Agent: Organizing the Data R&D Process

This solution does not merely add a dialogue entry to an existing platform; it reorganizes scattered capabilities in the data R&D process. Query, verification, planning, preview, submission, and status check are first encapsulated as tools with clear parameters and stable returns. Data sources, transformation relationships, and delivery targets are expressed uniformly through a domain DSL. Forge Agent then invokes these capabilities according to user needs, chaining the complete R&D flow.

Layer responsibilities:

R&D Entry – Forge Agent : Requirement clarification, task decomposition, Skill routing, tool selection, result explanation, and human handoff.

Tool Interface – Sindri Plugin/MCP : Encapsulates context query, DSL construction, verification, planning, submission, and status query as stable tools.

Domain Backend – Sindri Backend : Parses domain semantics, binds platform objects, generates execution plans, builds requests, and submits them.

Execution Base – ADF : Creates DAGs, organizes node dependencies, schedules runs, and maintains task and instance states.

To make this flow stable for production, Forge Agent implements the following mechanisms:

Skill : Precipitates operation flows and domain knowledge for specific scenarios.

Plugin/MCP : Executes standard operations such as compilation, verification, generation, submission, and query.

Hook : Performs permission checks, auditing, and human confirmation at critical nodes like submission.

Trace/Trajectory : Records actual invocation processes for problem localization and retrospectives.

Golden Case/Eval : Verifies via fixed cases whether version changes affect existing capabilities.

With this division of labor, the model only handles requirement understanding and flow organization, while production operations remain constrained by platform interfaces, permission rules, and execution results. This reduces manual hopping across systems while preserving the determinism of existing platforms in planning, scheduling, permissions, and auditing.

Sindri: From Domain Semantics to Trusted Execution Plans

Forge Agent organizes the R&D process for humans; Sindri explains to the platform how a data processing job should actually execute. The two interact via deterministic tools, not by letting the model directly assemble ADF requests.

A typical chain can be summarized as:

Requirement clarification → Context acquisition → DSL construction → Structural validation → Sindri authoritative planning → Submission preview → Explicit confirmation → Sindri submission → ADF lifecycle → Status feedback.

4.1 DSL Expresses "What to Do"

The domain DSL uses stable semantics to describe data objects, transformation operators, filtering, joins, aggregation, UDF/UDTF, and output targets. It does not concern itself with any specific engine or resource configuration.

Below is a simplified core topology of the Chaoxing Index Pipeline. This task simultaneously processes real-time product changes, real-time price changes, and offline full product data, forming real-time and offline output branches:

def build_dsl():
    ctx = SemanticContext(
        name="chaoxing_hk_pipeline",
        owner="<owner>"
    )

    # Two real-time change streams and one offline full data; specific SQL maintained separately
    sku_changes = ctx.sql(build_sku_change_stream_sql())
    price_changes = ctx.sql(build_price_change_stream_sql())
    full_skus = ctx.sql(build_full_sku_batch_sql())

    # Real-time branch: merge product and price changes, then merge with wide table state in HBase
    realtime = (
        ctx.union(sku_changes, price_changes)
        .merge("wide_table_chaoxing_hk")
        .on("rowkey")
        .alias("realtime_sku")
    )

    # Join main site product data to supplement extended attributes needed for index
    main_skus = ctx.from_table("hbase_main_search_data").alias("main_sku")
    enriched = (
        realtime.left_join(main_skus)
        .on("realtime_sku.sku_id" = main_sku.rowkey")
        .columns("sku_union_expand_ids")
    )

    # Offline branch: full data merged with historical wide table
    full_snapshot = (
        full_skus.merge("wide_table_chaoxing_hk")
        .on("rkey")
        .alias("full_snapshot")
    )

    # Two terminal branches deliver real-time messages and offline index data respectively
    enriched.to("jdq_search_index_event")
    full_snapshot.to("search_index_chaoxing_snapshot")
    return ctx

This example not only connects Source and Sink but also expresses multi-source input, union, state merge, cross-domain left_join, and multi-branch output within a SemanticContext. The DSL only describes data dependencies and business topology: which Flink job type the real-time branch falls into, how the offline branch is split, which operators and resources are needed — all generated by the Sindri backend based on table types and runtime configurations.

4.2 Improving DSL Generation Accuracy

The most prominent problem in DSL development today is not "can code be generated" but whether the output is stable and conforms to the current SDK and backend rules.

Incomplete requirement information : When source tables, field mappings, output targets, or run modes are unconfirmed, the model tends to auto-fill business parameters.

Overlapping Skill responsibilities : Requirement extraction, DSL development, fault diagnosis, and submission intervening simultaneously leads to irrelevant answers or premature submission during development.

Relying on outdated memory : After SDK interfaces, table types, and backend planning rules change, syntax may look correct but topology or execution shape no longer matches the current implementation.

Sindri Plugin does not rely solely on prompt constraints; it breaks the development flow into scenario routing, requirement contracts, template selection, and progressive validation. Responsibilities:

Local validation only proves the DSL structure is basically sound; final task splitting, engine selection, and resource planning are subject to backend optimization. Submission also uses a two-phase gate: first generate a preview, only execute real submission after explicit user confirmation. Through the chain of "single-scenario routing → requirement contract → standard examples → SDK facts → local validation → backend planning → human confirmation", DSL accuracy is transformed from a one-shot generation problem into an inspectable, rollback-capable engineering process.

ADF: Carrying DAG, Scheduling, State, and Lifecycle

ADF serves as a general execution and control base. It receives materialized task requests from the domain backend, organizes DAG nodes and dependencies, manages scheduling, execution, release, and status query, so upper-layer domain semantics no longer need to rebuild these capabilities.

ADF supports Spark, Flink, Shell, JDOS, and other task forms, continuously extending via an adaptation mechanism. The domain backend only describes execution intent; the specific engine on which it lands is decided by the adaptation layer, thereby isolating engine diversity within ADF.

Unified Base Pool: Stable Semantics Connecting Diverse Execution and Delivery

The unified base pool does not rewrite all tasks into the same code; instead it establishes a stable middle layer: upstream sources can change, downstream products can change, execution forms can change, but domain semantics, planning interfaces, submission boundaries, and state models remain consistent.

Quality, Lineage, and Self-Healing: Evidence First, Automation Later

Quality checks, reconciliation, lineage, monitoring, and recovery are indispensable modules for trustworthy data production. The data quality system currently focuses on "making quality assurance solid" and consists of four capabilities:

Data Reconciliation : Full and sampling verification, difference attribution, confirming that multiple outputs are truly consistent in definition.

Data Quality : Rule validation, completeness and consistency checks, threshold alerts, turning "whether data meets standards" into measurable constraints.

Lineage Tracking : Field-level lineage, impact analysis, and traceability, letting every product answer "where from, who affected".

Monitoring Dashboard : Chain health, SLA and timeliness, anomaly observation, presenting runtime trustworthiness in real time.

The data quality system's current positioning is quality assurance — solidifying these four capabilities. The goal is to further evolve into a trustworthy base that runs through the entire chain: precipitating evidence from quality capabilities and gradually opening it for automation to consume.

A judgment on the planning direction: the self-healing loop should not target "full automation" from the start; instead, evidence should be connected first, then automation expanded within the scope supported by evidence. Reconciliation, quality, lineage, and monitoring produce exactly the evidence foundation for this direction.

Along this direction:

Planning Evidence : Record target environment, semantic structure, execution nodes, resource plans, and confirmation results.

Submission Evidence : Record preview, permission judgment, submission results, and Sindri/ADF object identifiers.

Runtime Evidence : Return DAG, tasks, instances, logs, and structured errors, allowing the Agent to continue handling based on facts.

On these foundations, local validation, diagnostic suggestions, and limited retries can be gradually integrated. True self-healing still requires stable state machines, idempotency mechanisms, risk grading, retry limits, result re-verification, and human takeover — which is why "evidence first, automation later" is safer: the more complete the evidence chain, the more actions can be safely handed to automation.

Partial Effect Demonstration

Task production only needs a description document or a prompt to generate a base pool task; the corresponding execution chain is automatically generated, and key data is strictly guarded during execution to avoid user errors.

After task generation, corresponding execution files are produced. Upon user submission, resource archiving is synchronized for subsequent iterative development and overall system regression testing.

In the data quality system, base pool task lineage is automatically drawn (currently synchronously drawing other tasks associated with each individual step), and clicking a step synchronously displays detailed information for that step, facilitating full-chain tracking and troubleshooting.

Conclusion

Data Nirvana is not about adding a chat entry to the old chain, but reshaping the production method from business intent to data value.

Forge Agent, as the unified AI facade, organizes requirements, knowledge, processes, and feedback; Sindri Plugin/MCP turns deterministic capabilities into discoverable, composable tools; Sindri backend guards the planning and submission boundaries; ADF handles DAG, scheduling, state, and lifecycle. Together they form a collaboration chain of "intelligent organization, deterministic planning, controlled submission, stable execution".

When responsibility boundaries are clearly drawn, when every planning has basis, every submission has confirmation, every run has state, AI becomes not just a content generator but a trustworthy data productive force.

Original Source

Signed-in readers can open the original source through BestHub's protected redirect.

Sign in to view source
Republication Notice

This article has been distilled and summarized from source material, then republished for learning and reference. If you believe it infringes your rights, please contactadmin@besthub.devand we will review it promptly.

Data PipelineDSLData QualityData PlatformAI AgentADFForge AgentSindri
JD Tech Talk
Written by

JD Tech Talk

Official JD Tech public account delivering best practices and technology innovation.

0 followers
Reader feedback

How this landed with the community

Sign in to like

Rate this article

Was this worth your time?

Sign in to rate
Discussion

0 Comments

Thoughtful readers leave field notes, pushback, and hard-won operational detail here.