Skip to main content

Some fun with Apache Wayang and Spark / Tensorflow

Struggling with delivery, architecture alignment, or platform stability?

I help teams fix systemic engineering issues: processes, architecture, and clarity.
→ See how I work with teams.


Apache Wayang is an open-source Federated Learning (FL) framework developed by the Apache Software Foundation. It provides a platform for distributed machine learning, with a focus on ease of use and flexibility. It supports multiple FL scenarios and provides a variety of tools and components for building FL systems. It also includes support for various communication protocols and data formats, as well as integration with other Apache projects such as Apache Kafka and Apache Pulsar for data streaming. The project aims to make it easier to develop and deploy machine learning models in decentralized environments.

It's important to note that this are just examples and they may not be the way for your project to interact with Apache Wayang, you may need to check the documentation of the Apache Wayang project (https://wayang.apache.org) to see how to interact with it. I just point out how easy it is to use different languages to interact between Wayang and Spark.

Also, you need to make sure that you have the correct permissions and credentials to interact with the Wayang API and make changes to the Spark cluster.

Wayang - Scala - Spark:

import org.apache.wayang.{Wayang, WayangClient}

class SparkScaler(wayangUrl: String) {
    val wayang = new WayangClient(wayangUrl)

    def scaleUp(numWorkers: Int): Unit = {
        wayang.addWorkers(numWorkers)
    }

    def scaleDown(numWorkers: Int): Unit = {
        wayang.removeWorkers(numWorkers)
    }
}

The SparkScaler class takes a single parameter, the URL of the Wayang API endpoint, when it is initialized. The scaleUp() method can be called to add a specified number of workers to the Spark cluster, and the scaleDown() method can be called to remove a specified number of workers.

Wayang - Python - Spark

from apache_wayang import Wayang

class SparkScaler:
    def __init__(self, wayang_url):
        self.wayang = Wayang(wayang_url)

    def scale_up(self, num_workers):
        self.wayang.add_workers(num_workers)

    def scale_down(self, num_workers):
        self.wayang.remove_workers(num_workers)

The SparkScaler class takes a single parameter, the URL of the Wayang API endpoint, when it is initialized. The scale_up() method can be called to add a specified number of workers to the Spark cluster, and the scale_down() method can be called to remove a specified number of workers.

Wayang - Java Streams - Spark

import org.apache.wayang.WayangClient;
import java.util.stream.IntStream;

public class SparkScaler {
    private WayangClient wayang;

    public SparkScaler(String wayangUrl) {
        wayang = new WayangClient(wayangUrl);
    }

    public void scaleUp(int numWorkers) {
        IntStream.range(0, numWorkers).forEach(i -> wayang.addWorker());
    }

    public void scaleDown(int numWorkers) {
        IntStream.range(0, numWorkers).forEach(i -> wayang.removeWorker());
    }
}

The SparkScaler class takes a single parameter, the URL of the Wayang API endpoint, when it is initialized. The scaleUp() method can be called to add a specified number of workers to the Spark cluster, and the scaleDown() method can be called to remove a specified number of workers.

Iterate the K-Means clustering algorithm from Apache Wayang to TensorFlow

import org.apache.wayang.WayangClient;
import org.tensorflow.Graph;
import org.tensorflow.Session;
import org.tensorflow.Tensor;

public class KMeansIteration {
    private WayangClient wayang;
    private Graph graph;
    private Session session;

    public KMeansIteration(String wayangUrl, String modelPath) {
        wayang = new WayangClient(wayangUrl);
        graph = new Graph();
        graph.importGraphDef(modelPath);
        session = new Session(graph);
    }

    public void iterate(Tensor input) {
        Tensor wayangOutput = wayang.runKMeans(input);
        Tensor tfOutput = session.runner().feed("input", wayangOutput).fetch("output").run().get(0);
        // Perform further processing on tfOutput
    }
}

That's are only examples to show how easy it can be to get started with FL and also get involved into Wayang as a developer. Also consider to contribute to the project, check the project under wayang.apache.org 

The KMeansIteration class takes two parameters, the URL of the Wayang API endpoint and the path of the TensorFlow model, when it is initialized. The iterate() method can be called with an input Tensor, it will pass it to the Wayang's K-Means clustering algorithm, it will receive the output, and then will pass it to the TensorFlow's model as an input.

If you need help with distributed systems, backend engineering, or data platforms, check my Services.

Most read articles

Building a Model-Agnostic Multi-Agent System with OpenClaw

Over one week we rebuilt our AI stack around OpenClaw’s multi-agent architecture to avoid provider lock-in and stop wasting premium tokens. By aligning models to tasks, diversifying fallbacks across providers, enforcing minimal tool access, and switching to memory-first workflows with ephemeral sessions, we reduced token usage per task by about 70% and cut our monthly bill by 77% while improving operational resilience. How We Achieved 77% Cost Reduction and Provider Independence Over the past week, we rebuilt our AI infrastructure around OpenClaw’s multi-agent architecture. The result was a 77% cost reduction , provider independence , and a delegation system that routes work to the most cost-effective model for each job. Below is the technical journey of optimizing a 7-agent squad with OpenClaw. The Challenge: Model Provider Lock-In We started with a simple problem: our entire squad defaulted to a single model provider. This created three issues: Cost inefficiency beca...

BacNet => MQTT in Production: The Real Cost of Bridging BACnet to MQTT at Scale

bacnet2mqtt looks simple in a README and expensive in production. Once BACnet polling, reconnection behavior, stale state, and MQTT publishing collide, teams discover they are not deploying a lightweight adapter but operating infrastructure. This article breaks down where bacnet2mqtt works, where it becomes a bottleneck, and which production patterns reduce the operational damage before incidents, backlogs, and silent data loss turn a building integration into a long-running engineering problem. I inherited a building controls integration problem 18 months ago. Three office floors. 217 BACnet sensors covering temperature, occupancy, and HVAC actuators. The data was trapped inside the building automation network while the business wanted analytics, reporting, and compliance visibility in the data platform. The obvious answer looked easy enough: deploy bacnet2mqtt, bridge BACnet into MQTT, and push the stream into the lakehouse stack. The repository made it sound like a w...

Get Apache Flume 1.3.x running on Windows

Since we found an increasing interest in the flume community to get Apache Flume running on Windows systems again, I spent some time to figure out how we can reach that. Finally, the good news - Apache Flume runs on Windows. You need some tweaks to get them running. Prerequisites Build system: maven 3x, git, jdk1.6.x, WinRAR (or similar program) Apache Flume agent: jdk1.6.x, WinRAR (or similar program), Ultraedit++ or similar texteditor Tweak the Windows build box 1. Download and install JDK 1.6x from Oracle 2. Set the environment variables    => Start - type " env " into the search box, select " E dit system environment variables ", click Environment Variables, Select " New " from the " Systems variables " box, type " JAVA_HOME " into " variable name " and the path to your JDK installation into "Variable value" (Example:  C:\Program Files (x86)\Java\jdk1.6.0_33 ) 3. Download maven from Apache 4. Set...