This file provides guidance to Claude Code (claude.ai/code) when working with code in this repository.
Mercari Pipeline is a configuration-driven data pipeline framework built on Apache Beam. It enables running various data pipelines on Cloud Dataflow, Apache Spark, and Apache Flink without writing code - just YAML/JSON configuration files. It also includes a Server feature as an auxiliary tool to create, debug, and deploy pipelines.
module/- Pipeline module configoptions/- Pipeline options configsystem.md- Pipeline system configserver/- Pipeline server
docs/developer/pipeline/pipeline/- Document for the Pipeline itselfdocs/developer/pipeline/server/- Document for Pipeline Server
MPipeline.java- Main entry point that loads config, creates Pipeline, and applies sources/transforms/sinks
Config.java- Parses YAML/JSON config from GCS, local files, base64, or Parameter Manager- Supports FreeMarker templating with
${args.varName}syntax for runtime variable substitution - Config can import other config files via
system.imports
The pipeline is built from three module types, discovered via classpath scanning with @Module annotations:
Sources (module/source/) - Data input:
- BigQuery, Spanner, Bigtable, Datastore, Firestore
- PubSub, Kafka (streaming)
- JDBC, Storage (files), HTTP, Drive
Transforms (module/transform/) - Data processing:
- SelectTransform - Field selection/transformation with expressions
- AggregationTransform - Grouping and aggregation
- BeamSQLTransform - SQL-based processing
- OnnxTransform - ML inference
- PartitionTransform, LookupTransform, etc.
Sinks (module/sink/) - Data output:
- BigQuery, Spanner, Bigtable, Datastore, Firestore
- PubSub, Storage (files), JDBC, Iceberg
Module.java- Base class for all modulesMCollection.java- Wrapper around PCollection with schema metadataMCollectionTuple.java- Container for multiple named MCollectionsMElement.java- Universal data element handling Avro, Row, Entity, Struct typesSchema.java- Unified schema representation
schema/- Schema utilities for Avro, Row, Entity, Struct, Proto conversionspipeline/- Pipeline utilities (Filter, Select, Aggregation, Query)cloud/- Cloud service utilities (GCS, BigQuery, Spanner, PubSub)domain/- Domain-specific utilities (ML, text analysis, JDBC)
PipelineApiServer.java- REST API server for pipeline validationPipelineMcpStreamableServer.java- MCP server for AI integration
- Create class in appropriate package (
source/,transform/, orsink/) - Extend
Source,Transform, orSinkbase class - Add
@Module(name = “modulename”)annotation - Implement
expand()method returningMCollectionTuple - Module is auto-discovered via classpath scanning
system:
args:
myVar: “value”
imports:
- base: “gs://bucket/”
files: [“common.yaml”]
sources:
- name: input1
module: bigquery
parameters:
query: “SELECT * FROM table”
transforms:
- name: process1
module: select
inputs: [input1]
parameters:
fields: [...]
sinks:
- name: output1
module: spanner
inputs: [process1]
parameters:
projectId: myproject
instanceId: myinstance
databaseId: mydatabase
table: mytable# Build and create FlexTemplate container (default: Dataflow runner)
mvn clean package -DskipTests -Dimage={region}-docker.pkg.dev/{project}/{repo}/dataflow:latest
# Build for local execution (DirectRunner)
mvn clean package -DskipTests -Pdirect -Dimage=“{region}-docker.pkg.dev/{project}/{repo}/direct”
# Build API server
mvn clean package -DskipTests -Pserver -Dimage=“{region}-docker.pkg.dev/{project}/{repo}/server”
# Run all tests
mvn test
# Run a single test class
mvn test -Dtest=ConfigTest
# Run a single test method
mvn test -Dtest=ConfigTest#testMethodNamedataflow(default) - Cloud Dataflow runnerdirect- Local DirectRunner for testingserver- Pipeline API server (WAR)prism- PrismRunner for portable executionflink- Apache Flink runnerspark- Apache Spark runner
- Java 21
- Apache Beam 2.70.0
- Google Cloud Platform SDKs (BigQuery, Spanner, Datastore, etc.)
- Jetty EE11 12 (For Server)