Created
March 1, 2016 07:36
-
-
Save forthxu/cfd1fdf31a361b005eed to your computer and use it in GitHub Desktop.
zeromq php的使用,需要装php扩展
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| <?php | |
| /* | |
| * Hello World client | |
| * Connects REQ socket to tcp://localhost:5555 | |
| * Sends "Hello" to server, expects "World" back | |
| * @author Ian Barber <ian (dot) barber (at) gmail (dot) com> | |
| */ | |
| $context = new ZMQContext (); | |
| // Socket to talk to server | |
| echo "Connecting to hello world server...\n"; | |
| $socket = new ZMQSocket ($context, ZMQ::SOCKET_REQ); | |
| $socket->connect ("tcp://localhost:5555"); | |
| for($request_nbr = 0; $request_nbr != 10; $request_nbr++) { | |
| printf ("Sending request %d...\n", $request_nbr); | |
| $socket->send ("Hello"); | |
| $reply = $socket->recv (); | |
| printf ("Received reply %d: [%s]\n", $request_nbr, $reply); | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| <?php | |
| /* | |
| * Simple request-reply broker | |
| * @author Ian Barber <ian(dot)barber(at)gmail(dot)com> | |
| */ | |
| // Prepare our context and sockets | |
| $context = new ZMQContext(); | |
| $frontend = new ZMQSocket($context, ZMQ::SOCKET_ROUTER); | |
| $backend = new ZMQSocket($context, ZMQ::SOCKET_DEALER); | |
| $frontend->bind("tcp://*:5559"); | |
| $backend->bind("tcp://*:5560"); | |
| // Initialize poll set | |
| $poll = new ZMQPoll(); | |
| $poll->add($frontend, ZMQ::POLL_IN); | |
| $poll->add($backend, ZMQ::POLL_IN); | |
| $readable = $writeable = array(); | |
| // Switch messages between sockets | |
| while(true) { | |
| $events = $poll->poll($readable, $writeable); | |
| foreach($readable as $socket) { | |
| if($socket === $frontend) { | |
| // Process all parts of the message | |
| while(true) { | |
| $message = $socket->recv(); | |
| // Multipart detection | |
| $more = $socket->getSockOpt(ZMQ::SOCKOPT_RCVMORE); | |
| $backend->send($message, $more ? ZMQ::MODE_SNDMORE : null); | |
| if(!$more) { | |
| break; // Last message part | |
| } | |
| } | |
| }else if($socket === $backend) { | |
| $message = $socket->recv(); | |
| // Multipart detection | |
| $more = $socket->getSockOpt(ZMQ::SOCKOPT_RCVMORE); | |
| $frontend->send($message, $more ? ZMQ::MODE_SNDMORE : null); | |
| if(!$more) { | |
| break; // Last message part | |
| } | |
| } | |
| } | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| <?php | |
| /* | |
| * Weather update server | |
| * Binds PUB socket to tcp://*:5556 | |
| * Publishes random weather updates | |
| * @author Ian Barber <ian(dot)barber(at)gmail(dot)com> | |
| */ | |
| // Prepare our context and publisher | |
| $context = new ZMQContext(); | |
| $publisherSocket = $context->getSocket(ZMQ::SOCKET_PUB); | |
| $publisherSocket->bind("tcp://*:5556"); | |
| while (true) { | |
| // Get values that will fool the boss | |
| $zipcode = mt_rand(0, 100000); | |
| $temperature = mt_rand(-80, 135); | |
| $relhumidity = mt_rand(10, 60); | |
| // Send message to all subscribers | |
| $update = sprintf ("%05d %d %d", $zipcode, $temperature, $relhumidity); | |
| $publisherSocket->send($update); | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| <?php | |
| /* | |
| * Task sink | |
| * Binds PULL socket to tcp://localhost:5558 | |
| * Collects results from workers via that socket | |
| * @author Ian Barber <ian(dot)barber(at)gmail(dot)com> | |
| */ | |
| // Prepare our context and socket | |
| $context = new ZMQContext(); | |
| $receiver = new ZMQSocket($context, ZMQ::SOCKET_PULL); | |
| $receiver->bind("tcp://*:5558"); | |
| // Wait for start of batch | |
| $string = $receiver->recv(); | |
| // Start our clock now | |
| $tstart = microtime(true); | |
| // Process 100 confirmations | |
| $total_msec = 0; // Total calculated cost in msecs | |
| for ($task_nbr = 0; $task_nbr < 100; $task_nbr++) { | |
| $string = $receiver->recv(); | |
| if($task_nbr % 10 == 0) { | |
| echo ":"; | |
| } else { | |
| echo "."; | |
| } | |
| } | |
| $tend = microtime(true); | |
| $total_msec = ($tend - $tstart) * 1000; | |
| echo PHP_EOL; | |
| printf ("Total elapsed time: %d msec", $total_msec); | |
| echo PHP_EOL; |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| <?php | |
| /* | |
| * Task ventilator | |
| * Binds PUSH socket to tcp://localhost:5557 | |
| * Sends batch of tasks to workers via that socket | |
| * @author Ian Barber <ian(dot)barber(at)gmail(dot)com> | |
| */ | |
| $context = new ZMQContext(); | |
| // Socket to send messages on | |
| $sender = new ZMQSocket($context, ZMQ::SOCKET_PUSH); | |
| $sender->bind("tcp://*:5557"); | |
| echo "Press Enter when the workers are ready: "; | |
| $fp = fopen('php://stdin', 'r'); | |
| $line = fgets($fp, 512); | |
| fclose($fp); | |
| echo "Sending tasks to workers…", PHP_EOL; | |
| // The first message is "0" and signals start of batch | |
| $sender->send(0); | |
| // Send 100 tasks | |
| $total_msec = 0; // Total expected cost in msecs | |
| for ($task_nbr = 0; $task_nbr < 100; $task_nbr++) { | |
| // Random workload from 1 to 100msecs | |
| $workload = mt_rand(1, 100); | |
| $total_msec += $workload; | |
| $sender->send($workload); | |
| } | |
| printf ("Total expected cost: %d msec\n", $total_msec); | |
| sleep (1); // Give 0MQ time to deliver |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| <?php | |
| /* | |
| * Hello World server | |
| * Binds REP socket to tcp://*:5555 | |
| * Expects "Hello" from client, replies with "World" | |
| * @author Ian Barber <ian (dot) barber (at) gmail (dot) com> | |
| */ | |
| $context = new ZMQContext (1); | |
| // Socket to talk to clients | |
| $socket = new ZMQSocket ($context, ZMQ::SOCKET_REP); | |
| $socket->bind ("tcp://*:5555"); | |
| while(true) { | |
| // Wait for next request from client | |
| $request = $socket->recv (); | |
| printf ("Received request: [%s]\n", $request); | |
| // Do some 'work' | |
| //sleep (1); | |
| // Send reply back to client | |
| $socket->send ("World"); | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| <?php | |
| /* | |
| * Weather update client | |
| * Connects SUB socket to tcp://localhost:5556 | |
| * Collects weather updates and finds avg temp in zipcode | |
| * @author Ian Barber <ian(dot)barber(at)gmail(dot)com> | |
| */ | |
| $context = new ZMQContext(); | |
| // Socket to talk to server | |
| echo "Collecting updates from weather server…", PHP_EOL; | |
| $subscriberSocket = new ZMQSocket($context, ZMQ::SOCKET_SUB); | |
| $subscriberSocket->connect("tcp://localhost:5556"); | |
| // Subscribe to zipcode, default is NYC, 10001 | |
| $filter = $_SERVER['argc'] > 1 ? $_SERVER['argv'][1] : "10001"; | |
| $subscriberSocket->setSockOpt(ZMQ::SOCKOPT_SUBSCRIBE, $filter); | |
| // Process 100 updates | |
| $total_temp = 0; | |
| for ($update_nbr = 0; $update_nbr < 100; $update_nbr++) { | |
| $string = $subscriberSocket->recv(); | |
| echo $string."\n"; | |
| sscanf ($string, "%d %d %d", $zipcode, $temperature, $relhumidity); | |
| $total_temp += $temperature; | |
| } | |
| printf ("Average temperature for zipcode '%s' was %dF\n",$filter, (int) ($total_temp / $update_nbr)); |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| <?php | |
| /* | |
| * Task worker | |
| * Connects PULL socket to tcp://localhost:5557 | |
| * Collects workloads from ventilator via that socket | |
| * Connects PUSH socket to tcp://localhost:5558 | |
| * Sends results to sink via that socket | |
| * @author Ian Barber <ian(dot)barber(at)gmail(dot)com> | |
| */ | |
| $context = new ZMQContext(); | |
| // Socket to receive messages on | |
| $receiver = new ZMQSocket($context, ZMQ::SOCKET_PULL); | |
| $receiver->connect("tcp://localhost:5557"); | |
| // Socket to send messages to | |
| $sender = new ZMQSocket($context, ZMQ::SOCKET_PUSH); | |
| $sender->connect("tcp://localhost:5558"); | |
| // Process tasks forever | |
| while (true) { | |
| $string = $receiver->recv(); | |
| // Simple progress indicator for the viewer | |
| echo $string, PHP_EOL; | |
| // Do the work | |
| usleep($string * 1000); | |
| // Send results to sink | |
| $sender->send(""); | |
| } |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment