Product
How EdgeMQ Works
From HTTP requests to queryable data in your lake
EdgeMQ is a managed HTTP → S3 ingest layer. Producers send NDJSON over HTTPS to a regional /ingest endpoint. EdgeMQ:
- Accepts and authenticates the request at the edge.
- Writes it to a durable write-ahead log (WAL) on NVMe.
- Groups WAL entries into segments.
- Materializes each sealed segment into segments and/or Parquet (raw or views) in your S3 bucket and prefix.
- Writes a commit marker that lists the artifacts for that segment so your data stack knows what's safe to read.
From there, your lakehouse and warehouse tools-Snowflake, Databricks, ClickHouse, DuckDB, Postgres-treat S3 as a continuously updated Bronze layer.
This page walks through that flow and explains the guarantees you get.
High-level architecture
Think of EdgeMQ as a thin but very opinionated layer between producers and your S3 lake:
Producers
Services, devices, jobs send NDJSON over HTTPS
EdgeMQ
Edge nodes: TLS, WAL, queues, compression
Your S3
Compressed segments and/or Parquet files + commit markers
Consumers
Snowflake, Databricks, ClickHouse, DuckDB, Postgres
Producers
Services, devices, jobs, and partners send NDJSON over HTTPS.
EdgeMQ edge nodes
Regional microVMs that terminate TLS, append to WAL, manage queues, and ship segments to S3.
Your S3 bucket
The durable, long-term store for segment files, Parquet (raw or schema-aware views), and commit markers.
Consumers
Your existing tools: Snowflake, Databricks, ClickHouse, DuckDB, Postgres, and more.
EdgeMQ owns the ingest complexity; your stack continues to own storage, processing, and analytics.
Ingest at the edge: /ingest over HTTPS
Simple HTTP contract
Every producer talks to EdgeMQ via a regional HTTPS endpoint:
curl -X POST "https://<region>.edge.mq/ingest" \
-H "Authorization: Bearer $EDGEMQ_TOKEN" \
-H "Content-Type: application/x-ndjson" \
--data-binary @events.ndjson- ▸Transport: HTTPS, terminated close to your producers.
- ▸Payload: NDJSON (application/x-ndjson) - one JSON object per line.
- ▸Auth: Token/JWT-based; configurable per environment.
The same pattern works for: backend services emitting business events, gateways forwarding device telemetry, batch jobs uploading bigger NDJSON dumps, partners sending webhooks (through a small adapter or directly).
Everyone gets the same answer: "POST NDJSON here; it will end up in S3."
Durability: WAL on NVMe
Once a request hits the edge, EdgeMQ's first job is not to lose it.
Per-tenant microVMs
Each account/region runs inside an isolated microVM with:
- ▸Its own NVMe-backed volume.
- ▸A single-writer write-ahead log (WAL).
- ▸Its own queues and network stack.
There is no shared WAL or disk between tenants.
Write → fsync → ack
For each request:
- EdgeMQ parses the NDJSON payload into frames.
- Appends those frames to the WAL.
- Syncs the WAL to disk (fsync or equivalent).
- Only then returns a success response to the client.
If the process or VM restarts, EdgeMQ: replays the WAL, resumes segment creation and S3 uploads, ensures acknowledged events are not lost.
Durability is local and explicit-not "we hope the upload finished".
Segments: batching for S3 and for your jobs
Continuous WAL writes are great for durability but not ideal for S3 or downstream jobs. EdgeMQ groups WAL entries into segments, then uses those sealed segments to produce the artifacts you read from S3 (compressed segment files and, when configured, Parquet files).
Open segment
Active segment accumulating new frames.
Seal
When a segment reaches a size or time threshold, EdgeMQ seals it (no more writes).
Compress
The sealed segment is compressed (e.g. zstd) to reduce bandwidth and storage.
Upload
The compressed segment is shipped to your S3 bucket via multipart upload.
This yields: reasonable object sizes in S3 (good for listing and throughput), natural time/size boundaries for downstream processing, efficient compression for NDJSON workloads.
When Parquet output is enabled for an endpoint, EdgeMQ produces Parquet from each sealed segment:
- ▸Uses the segment as the source of truth.
- ▸Writes one or more Parquet files into a partitioned S3 layout (for example, by tenant and date).
- ▸Choose a mode: raw/opaque (payload preserved as a field) or schema-aware views (typed columns generated from view definitions), depending on how you want to query.
Parquet: raw vs. views
Raw Parquet preserves the original record as a payload field alongside metadata columns. Schema-aware Parquet views produce typed columns so your tools can query tables directly.
Shipping to your S3 bucket
EdgeMQ doesn't store your data in its own buckets. It writes to your cloud account under your policies.
IAM-based access
For each environment you configure:
- ▸S3 bucket and prefix (for example,
s3://your-bucket/edge-events/prod/). - ▸An IAM Role that EdgeMQ can assume.
EdgeMQ:
- ▸Uses AWS STS (AssumeRole) to obtain short-lived credentials.
- ▸Is restricted by the IAM policy to: a specific bucket/prefix, specific S3 actions (e.g. PutObject, multipart upload operations, minimal reads).
You retain control over:
- ▸Encryption (KMS).
- ▸Lifecycle (retention, tiering).
- ▸Access policies and audits.
Multipart upload and retries
Segments are uploaded using multipart S3 uploads: robust to larger segment sizes, better resilience to transient network issues, easier to retry parts without re-sending the entire object.
Depending on how an endpoint is configured, EdgeMQ will upload:
- ▸Segment artifacts under a segments prefix (compressed WAL segments), and/or
- ▸Parquet artifacts under a Parquet prefix (partitioned by date)
All artifacts for a given segment share a common sequence number and are referenced from the same commit marker.
The WAL remains the source of truth until the upload and commit are complete.
Commit markers: knowing what's safe to read
Once a segment is successfully uploaded, EdgeMQ writes a commit marker to S3.
What is a commit marker?
A small JSON object that:
- ▸Lives under a predictable prefix (e.g.
commits/seg-XXXX.json). - ▸Acts as a manifest for that segment, listing the artifacts that exist (for example, one compressed segment file and several Parquet files) and their S3 keys.
- ▸Indicates that all required artifacts for this segment are fully and successfully stored.
Why they exist
Without commit markers, consumers have to guess: "Is this file complete or still being uploaded?" "Should I retry if the file looks truncated?"
With commit markers:
- ▸Consumers only process segments with a corresponding marker.
- ▸Simple, robust pattern: list commit markers, skip any markers already processed, and for each new marker read the referenced artifacts (segment file and/or Parquet files) and update your own checkpoint.
This provides a clean handoff between EdgeMQ and your pipelines.
Behavior under load: queues and backpressure
EdgeMQ is explicit about what happens when producers outpace capacity.
Bounded queues
Between the HTTP handler and the WAL/shipper, EdgeMQ uses bounded queues:
- ▸Buffers are limited in size.
- ▸When they're healthy, ingest latency is low and steady.
- ▸When they fill, EdgeMQ does not pretend everything is fine.
Backpressure via 503
When the system is saturated:
- ▸EdgeMQ returns HTTP 503 Service Unavailable.
- ▸It includes a Retry-After header to guide clients.
This is intentional: protects the WAL disk and CPU, prevents latency blow-ups for all tenants, makes overload behavior predictable for your client libraries.
If a request returns 200, it's durable; if it returns 503, the client knows it must retry.
Delivery semantics and guarantees
EdgeMQ is designed around clear, practical semantics.
At-least-once delivery
- ▸A request that receives a success response has been durably written to the WAL.
- ▸WAL replay on restart ensures data is shipped to S3.
- ▸In certain failure modes, a frame may be uploaded twice; it is not silently lost.
Downstream consumers should be idempotent or deduplicate when necessary (e.g. using natural keys in the payload).
Per-instance ordering
- ▸Within a single ingest instance: frames are written to the WAL in arrival order, that order is preserved within each segment.
- ▸Global total ordering across all regions/instances is not guaranteed; most architectures rely on: partitioning (by stream/tenant/source), timestamps and IDs embedded in events.
How your tools plug into EdgeMQ
EdgeMQ doesn't replace your warehouse or engine; it feeds them.
Snowflake
- ▸Use Snowpipe or scheduled
COPY INTO: - ▸Source EdgeMQ-managed S3 prefixes
- ▸Read either NDJSON segments or Parquet files, plus commit markers/manifest to know what's ready
- ▸Treat these prefixes as your raw / Bronze layer
Databricks / Spark
- ▸Hook up Autoloader or Spark batch jobs: input EdgeMQ Parquet prefixes (or segments if you prefer), using commit markers to drive incremental loads into Delta tables (Bronze/Silver/Gold).
- ▸Parquet output means you can treat EdgeMQ-managed prefixes as regular Parquet tables in your lakehouse.
ClickHouse / Postgres
- ▸Use EdgeMQ as a burst buffer: producers push to EdgeMQ → S3, loader jobs ingest into ClickHouse/Postgres at a controlled rate, decouple ingest spikes from database capacity.
DuckDB and notebooks
- ▸Analysts and data scientists query directly: use
read_json_autoover segment outputs, orread_parquetover Parquet outputs, depending on how the endpoint is configured. - ▸This gives you ad-hoc access to fresh data without standing up new services.
Today vs. roadmap
What works today
- ▸Ingest NDJSON over HTTPS via
/ingest. - ▸Durably log events to a per-tenant WAL.
- ▸Segment, compress, and upload to your S3 bucket.
- ▸Materialize sealed segments into compressed segment files and/or Parquet files, depending on endpoint configuration.
- ▸Emit commit markers (manifests) for each completed segment that list the artifacts and their S3 keys.
- ▸Integrate with warehouses/lakehouses using their existing file ingestion mechanisms (reading segments or Parquet (raw or views) directly from S3).
Where it's going
EdgeMQ is evolving into a richer, format-aware ingest layer:
- ▸Additional ingest formats (beyond NDJSON) over time.
- ▸Additional output formats on S3, building on segment + Parquet support:
- ▸CSV where needed
- ▸Table-friendly layouts (e.g. Iceberg/Delta-style directories) for lakehouse engines
- ▸Richer table formats on top of Parquet (e.g. Iceberg/Delta-style layouts)
The core idea stays the same: one HTTP → S3 ingest layer, multiple engines on top.
Recap
EdgeMQ gives you a thin but powerful building block:
For producers
A simple, global HTTPS endpoint that takes NDJSON and either durably accepts or clearly rejects with retry instructions.
For platform/infra
A standardized, isolated, IAM-integrated ingest plane with predictable overload behavior and clear observability.
For data & ML teams
A continuously updated S3 Bronze layer that feeds Snowflake, Databricks, ClickHouse, DuckDB, Postgres, and future engines-without bespoke ingest code.
What's next
If this is the kind of ingest layer you want in front of your lakehouse: