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.
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.
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_OLDESTto 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, useSTART_FROM_CURRENT. - The client ID argument helps you tell which client is logging which events in the console.
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>
Follow this order to avoid connection issues:
- Start the
HazelcastMapServerfirst—wait until it starts printing entry additions. - Launch the first client with:
java JetMapJournalClient 1 - 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.
- 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
capacityon 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_OLDESTinstead ofSTART_FROM_CURRENT.
内容的提问来源于stack exchange,提问作者Harshad Murtekar

