Prompt
I'm building something like Fivetran but with custom business logic, where do I start?
Latest observation
If you’re building “Fivetran + custom business logic,” start by not trying to build a full connector platform first. Start by designing the data movement + transformation loop for one or two sources end-to-end, then generalize.
1) Clarify what you’re actually building
Fivetran is mostly about:
- Extracting from SaaS APIs / databases
- Handling schema drift
- Incremental syncs
- Loading into a warehouse
- Reliability / retries / observability
Your “custom business logic” could mean:
- Transforming records during ingestion
- Enriching with internal data
- Deduping / filtering / routing data
- Running domain-specific sync rules
- Triggering actions based on source events
Decide early whether your product is:
- ELT pipeline platform with transformation hooks, or
- Operational data sync platform with business-rule workflows, or
- Reverse ETL / sync engine with custom routing and state.
That decision changes the architecture a lot.
2) Start with a narrow MVP
Pick:
- 1 source type: e.g. Postgres or Stripe
- 1 destination: e.g. Postgres, Snowflake, BigQuery
- 1 sync mode: incremental CDC or API polling
- 1 transform point: e.g. a per-record function or SQL transform
A good first MVP:
- Sync from Postgres
- Detect changes via CDC
- Apply a custom transform function
- Load into Snowflake/Postgres
- Support retries, checkpoints, and backfills
That gives you the hard parts without API throttling complexity.
3) Think in core system components
A basic architecture usually needs these parts:
A. Connector runtime
A standardized way to run connectors:
- Input: config, credentials, sync state
- Output: normalized records + schema events + checkpoints
- Must support: batching, retries, pagination, rate limits
B. State management
You need durable state for:
- Last synced cursor / LSN / timestamp
- Checkpoint per stream/table
- Schema versions
- Job run status
C. Normalization layer
Convert each source into a common model:
recordschemadeletion/tombstonecheckpointerror
This makes everything else easier.
D. Transformation/business logic layer
This is where your product differentiates. Examples:
- Row-level Python/JS hooks
- SQL transforms
- Rules engine
- Event-driven enrichments
- Conditional routing by tenant/source/table
E. Load layer
A loader that can:
- Batch inserts
- Merge/upsert
- Handle idempotency
- Support schema evolution
F. Orchestration and retries
You need:
- Job scheduler
- Worker queue
- Retry policy
- Dead-letter handling
- Backfill support
- Alerting/monitoring
4) Pick your execution model early
This is a big design choice.
Option 1: “Connector as code”
Users write Python/TypeScript code for custom logic.
- Pros: flexible, powerful
- Cons: sandboxing, security, support burden
Option 2: Declarative config + rules
Users define mapping/transforms in YAML/SQL/JSON.
- Pros: safer, easier to operate
- Cons: less flexible
Option 3: Hybrid
Built-in connectors + configurable transform steps + optional custom code sandbox. This is often the best long-term answer.
If you expect enterprise adoption, avoid arbitrary code execution in your core control plane unless you have a strong isolation model.
5) Build a minimal data model
You’ll want something like:
- Connection
- source type
- credentials
- destination
- sync config
- Sync Job
- start time
- end time
- status
- metrics
- Stream/Table
- schema
- cursor field
- checkpoint
- Transform
- config or code reference
- version
- Artifact
- logs, errors, batch stats
6) Use proven patterns
For the first version, copy the patterns that work:
For sources
- DB sources: use CDC if possible
- API sources: use pagination + cursor-based incremental sync
- Webhooks: treat as event ingestion, not authoritative state
For loads
- Land raw data first
- Transform after landing, or in a staging area
- Make writes idempotent
- Prefer upsert/merge over blind overwrite
For reliability
- At-least-once delivery
- Deduplicate downstream
- Make checkpoints durable before acknowledging completion
- Expect partial failures
7) Decide where business logic runs
This is the key product question.
Embed logic in ingestion
Good for:
- Filtering rows
- PII masking
- Basic enrichment
- Routing
Bad for:
- Complex joins
- Heavy computations
- Cross-source workflows
Run logic after landing
Good for:
- SQL transforms
- Auditable pipelines
- Reprocessing
Usually best if your “custom logic” is more than trivial.
A common design is:
- ingest raw data
- apply user-defined transformations
- publish curated tables/events
8) Don’t underestimate schema drift
This will bite you early. You need policies for:
- New columns
- Deleted columns
- Type changes
- Nullability changes
- Nested JSON changes
Decide whether you:
- Auto-evolve schema
- Block on breaking changes
- Quarantine bad records
- Version schemas per sync
9) Build observability from day one
You’ll need:
- Per-connector logs
- Record counts in/out
- Lag metrics
- Retry counts
- Source API rate limit metrics
- Freshness SLA tracking
- Error classification
If users can’t see why data is late or wrong, they won’t trust the product.
10) Suggested build order
Here’s a practical sequence:
Phase 1: Prototype
- One source, one destination
- Simple incremental sync
- Local runner
- Hardcoded transform hook
- Basic logging
Phase 2: Platform skeleton
- Job orchestration
- Durable state store
- Retry/backfill
- Schema tracking
- UI/API for connection config
Phase 3: Custom logic
- Transform DSL or sandboxed code
- Validation/testing for transforms
- Versioning of logic
- Preview mode on sample data
Phase 4: Scale and hardening
- Multi-tenant isolation
- Secrets management
- Connector SDK
- Observability/alerting
- Parallelization and rate limit handling
11) Good tech choices for a first version
A reasonable stack:
- Control plane: Node.js/TypeScript or Python
- Workers: Python is common for connectors
- Orchestration: Temporal, Celery, or a queue-based system
- State store: Postgres
- Queue: SQS/RabbitMQ/Kafka depending on scale
- Transforms: SQL or a sandboxed runtime
- Metadata/API: Postgres + REST/gRPC
- Observability: OpenTelemetry + Prometheus/Grafana
If you want stronger workflow guarantees, Temporal is a strong choice.
12) Biggest mistake to avoid
Don’t start by building:
- a connector marketplace
- a full UI
- dozens of integrations
- a custom low-code transformation engine
Start with:
- one source
- one destination
- one durable sync loop
- one place for business logic
That will tell you what your real product should be.
13) A very practical first milestone
Build this:
“Sync Stripe customers into a warehouse, apply a custom rule to classify customers, and keep it updated incrementally with retries and checkpoints.”
Why this works:
- API sync complexity
- business logic included
- stateful incremental behavior
- real operational value
If you want, I can help you with one of these next:
- a reference architecture diagram
- an MVP feature list
- a connector SDK design
- a data model / schema for sync jobs
- a 30/60/90-day build plan