Giant Eye Tech — An Eye in the Sky
Giant Eye TechAn Eye in the Sky
Back to all RFCs & Articles
IoT & Manufacturing

Real-Time Loom & Machine Floor Telemetry Using Apache Kafka and TimescaleDB

Bridging flaky industrial manufacturing Wi-Fi with edge SQLite buffering and sub-second TimescaleDB hypertable aggregation.

Giant Eye Tech Systems Practice(Industrial IoT Engineering Practice)
August 2026
12 min read

Data Loss Rate

0.000%

Zero lost pulses during Wi-Fi blackouts

Telemetry Volume

45,000 msgs/s

Aggregated across 1,200 loom sensors

Dashboard Latency

< 250ms

Real-time floor supervisor tablets

01The Engineering Bottleneck

In modern textile, solar, and manufacturing facilities, shop floors are hostile radio environments. Metal structural beams, high-power motor EMI, and warehouse walls cause intermittent 5 to 60-second Wi-Fi dropouts. Standard cloud-first HTTP polling drops packets, stalls production line operators, and corrupts shift productivity metrics.

02Architectural Design & Invariants

We engineer a Store-and-Forward edge topology. Industrial Linux terminals (Raspberry Pi CM4 or Siemens IPCs) hardwired to machine relays log sensor ticks and operator barcodes to an embedded WAL-enabled SQLite queue. A lightweight Go edge daemon batches and compresses telemetry, publishing to an Apache Kafka cluster whenever network connectivity is verified.

03Production Implementation Blueprint

go
Go Edge Agent Store-and-Forward Telemetry Loop
Strict Production Invariant
package main

import (
    "context"
    "database/sql"
    "time"
    "github.com/segmentio/kafka-go"
)

type TelemetryTick struct {
    MachineID string    `json:"machine_id"`
    LoomRPM   int       `json:"loom_rpm"`
    YarnPicks uint64    `json:"yarn_picks"`
    Timestamp time.Time `json:"timestamp"`
}

func syncWorker(ctx context.Context, localDb *sql.DB, kafkaWriter *kafka.Writer) {
    ticker := time.NewTicker(500 * time.Millisecond)
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            return
        case <-ticker.C:
            // Fetch uncommitted edge records
            rows, err := localDb.QueryContext(ctx, 
                "SELECT id, payload FROM pending_telemetry ORDER BY id ASC LIMIT 250")
            if err != nil {
                continue
            }
            
            var messages []kafka.Message
            var processedIds []int64

            for rows.Next() {
                var id int64
                var payload []byte
                if err := rows.Scan(&id, &payload); err == nil {
                    messages = append(messages, kafka.Message{
                        Key:   []byte("loom_telemetry"),
                        Value: payload,
                    })
                    processedIds = append(processedIds, id)
                }
            }
            rows.Close()

            if len(messages) == 0 {
                continue
            }

            // Publish batch to Kafka broker with 2s write deadline
            if err := kafkaWriter.WriteMessages(ctx, messages...); err == nil {
                // Delete processed batch from edge SQLite WAL
                markCommitted(localDb, processedIds)
            }
        }
    }
}

04Architectural Invariants & Rules of Thumb

  • Edge machines must operate in 100% disconnected offline mode indefinitely without blocking factory operators.
  • Use TimescaleDB hypertables with 7-day chunk intervals and continuous 1-minute rollup aggregates for instant mill dashboard rendering.
  • Kafka partition keying by machine ID preserves strict causal ordering across shift cycles.
  • Reduces shift reconciliation lag from 24 hours (paper clipboard tally) to under 800 milliseconds.
Enterprise Systems Practice

Facing similar architecture bottlenecks in your business?

We design and implement custom ERPs, high-throughput databases, and air-gapped private AI systems tailored for high-concurrency enterprise workloads.