You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Hazelcast Jet 0.5.1服务端-客户端示例运行问题咨询

Got it, let's walk through how to extend that RemoteMapJournalSource sample into a proper server-client setup with 1 server and 2 clients. I'll break this down into actionable steps with code snippets so you can follow along smoothly.

Step 1: Build the Server (With Event Journal Enabled Map)

First, we need a server that initializes Hazelcast with an Event Journal-enabled map, then populates 100 integer entries into it. Here's the full code:

import com.hazelcast.config.Config;
import com.hazelcast.core.Hazelcast;
import com.hazelcast.core.HazelcastInstance;
import com.hazelcast.core.IMap;

public class HazelcastMapServer {
    public static void main(String[] args) {
        // Configure Hazelcast with Event Journal for our target map
        Config serverConfig = new Config();
        serverConfig.getMapConfig("journaled-numbers")
                    .getEventJournalConfig()
                    .setEnabled(true)
                    .setCapacity(1000) // Buffer size for events; adjust based on your load
                    .setTimeToLiveSeconds(300); // How long events are retained
        
        // Start the Hazelcast server instance
        HazelcastInstance server = Hazelcast.newHazelcastInstance(serverConfig);
        IMap<Integer, Integer> numberMap = server.getMap("journaled-numbers");
        
        // Populate 100 integer entries into the map
        System.out.println("Server starting to populate map...");
        for (int i = 0; i < 100; i++) {
            numberMap.put(i, i * 2);
            System.out.printf("Server added entry: %d -> %d%n", i, i*2);
            // Add a small delay so clients can catch up in real-time
            try {
                Thread.sleep(75);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                break;
            }
        }
        
        System.out.println("Server finished populating map. Keeping instance running...");
        // Keep server alive indefinitely (you can add a shutdown hook if needed)
        Runtime.getRuntime().addShutdownHook(new Thread(server::shutdown));
    }
}

The key here is enabling the Event Journal on the journaled-numbers map—this is what lets clients stream the map's entry events later.

Step 2: Build the Client (Stream Map Journal Events)

Our two clients will connect to the server, then use Hazelcast Jet to consume the map's event journal. The client code is reusable for both instances; we'll pass a command-line argument to distinguish them in logs:

import com.hazelcast.client.config.ClientConfig;
import com.hazelcast.core.HazelcastClient;
import com.hazelcast.core.HazelcastInstance;
import com.hazelcast.jet.Jet;
import com.hazelcast.jet.JetInstance;
import com.hazelcast.jet.pipeline.Pipeline;
import com.hazelcast.jet.pipeline.Sinks;
import com.hazelcast.jet.pipeline.Sources;
import com.hazelcast.jet.pipeline.JournalInitialPosition;

public class JetMapJournalClient {
    public static void main(String[] args) {
        if (args.length == 0) {
            System.err.println("Please pass a client ID (e.g., 1 or 2) as an argument");
            System.exit(1);
        }
        String clientId = args[0];
        
        // Configure client to connect to the Hazelcast server
        ClientConfig clientConfig = new ClientConfig();
        clientConfig.getNetworkConfig().addAddress("127.0.0.1:5701"); // Point to your server's address
        
        // Initialize Hazelcast client and Jet instance
        HazelcastInstance hazelcastClient = HazelcastClient.newHazelcastInstance(clientConfig);
        JetInstance jetClient = Jet.newJetInstance(clientConfig);
        
        // Build the Jet pipeline to consume map journal events
        Pipeline pipeline = Pipeline.create();
        pipeline.readFrom(Sources.remoteMapJournal(
                "journaled-numbers",
                clientConfig,
                JournalInitialPosition.START_FROM_OLDEST, // Catch all existing entries + new ones
                event -> event // Pass the full journal event to the next stage
            ))
            .setName("remote-map-journal-source")
            .writeTo(Sinks.logger(event -> String.format(
                "Client %s received event: Key=%d, Value=%d, Event Type=%s",
                clientId,
                event.getKey(),
                event.getNewValue(),
                event.getType()
            )));
        
        // Submit the Jet job to start streaming
        jetClient.newJob(pipeline);
        System.out.printf("Client %s connected to server and started consuming events...%n", clientId);
        
        // Keep client alive
        Runtime.getRuntime().addShutdownHook(new Thread(() -> {
            jetClient.shutdown();
            hazelcastClient.shutdown();
        }));
    }
}

Quick Client Notes:

  • Use JournalInitialPosition.START_FROM_OLDEST to have clients pick up all 100 entries the server wrote (even if clients start after the server finishes populating). If you only want new entries after client startup, use START_FROM_CURRENT.
  • The client ID argument helps you tell which client is logging which events in the console.
Step 3: Dependencies (Critical for 0.5.1 Compatibility)

Make sure your pom.xml (if using Maven) has the correct compatible versions—Hazelcast Jet 0.5.1 pairs with Hazelcast 3.10.x:

<dependencies>
    <dependency>
        <groupId>com.hazelcast.jet</groupId>
        <artifactId>hazelcast-jet</artifactId>
        <version>0.5.1</version>
    </dependency>
    <dependency>
        <groupId>com.hazelcast</groupId>
        <artifactId>hazelcast-client</artifactId>
        <version>3.10.2</version>
    </dependency>
</dependencies>
Step 4: Run the Setup

Follow this order to avoid connection issues:

  1. Start the HazelcastMapServer first—wait until it starts printing entry additions.
  2. Launch the first client with: java JetMapJournalClient 1
  3. Launch the second client with: java JetMapJournalClient 2

You'll see both clients receive all 100 map entries in their logs, each clearly labeled with their client ID.

Troubleshooting Tips
  • If clients can't connect: Double-check the server's IP/port in the client config, and ensure no firewall is blocking port 5701.
  • If clients miss events: Increase the Event Journal's capacity on the server—if entries are written faster than the journal can buffer, older events get evicted.
  • If clients only get new entries: Verify you're using START_FROM_OLDEST instead of START_FROM_CURRENT.

内容的提问来源于stack exchange,提问作者Harshad Murtekar

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.22 09:14:23