-
Notifications
You must be signed in to change notification settings - Fork 1.6k
Expand file tree
/
Copy pathfileio3.java
More file actions
134 lines (118 loc) · 4.14 KB
/
Copy pathfileio3.java
File metadata and controls
134 lines (118 loc) · 4.14 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
import org.zeromq.ZContext;
import org.zeromq.ZFrame;
import org.zeromq.ZMQ;
import org.zeromq.ZMQ.Socket;
import org.zeromq.ZThread;
import java.io.File;
import java.io.FileInputStream;
import java.io.FileNotFoundException;
import java.io.IOException;
import java.util.Arrays;
// File Transfer model #3
//
// In which the client requests each chunk individually, using
// command pipelining to give us a credit-based flow control.
public class Fileio3 {
private static final int PIPELINE = 10;
private static final int CHUNK_SIZE = 250000;
// The main task starts the client and server threads; it's easier
// to test this as a single process with threads, than as multiple
// processes:
public static void main(String[] args) {
ZContext ctx = new ZContext();
// Start child threads
ZThread.fork(ctx, new Fileio1.Server());
ZMQ.Socket client = ZThread.fork(ctx, new Fileio1.Client());
// Loop until client tells us it's done
client.recvStr();
// Kill server thread
ctx.destroy();
}
static class Client implements ZThread.IAttachedRunnable {
@Override
public void run(Object[] args, ZContext ctx, Socket pipe) {
Socket dealer = ctx.createSocket(ZMQ.DEALER);
dealer.connect("tcp://127.0.0.1:6000");
// Up to this many chunks in transit
int credit = PIPELINE;
int total = 0; // Total bytes received
int chunks = 0; // Total chunks received
int offset = 0; // Offset of next chunk request
while (true) {
while (credit > 0) {
// Ask for next chunk
dealer.sendMore("fetch");
dealer.sendMore(String.valueOf(offset));
dealer.send(String.valueOf(CHUNK_SIZE));
offset += CHUNK_SIZE;
credit--;
}
ZFrame chunk = ZFrame.recvFrame(dealer);
if (chunk.getData() == null)
break; // Shutting down, quit
chunks++;
credit++;
int size = chunk.size();
chunk.destroy();
total += size;
if (size < CHUNK_SIZE)
break; // Last chunk received; exit
}
System.out.printf("%d chunks received, %d bytes\n", chunks, total);
pipe.send("OK");
}
}
// The rest of the code is exactly the same as in model 2, except
// that we set the HWM on the server's ROUTER socket to PIPELINE
// to act as a sanity check.
// The server thread waits for a chunk request from a client,
// reads that chunk and sends it back to the client:
static class Server implements ZThread.IAttachedRunnable {
@Override
public void run(Object[] args, ZContext ctx, Socket pipe) {
File file = new File("testdata");
FileInputStream fr;
try {
fr = new FileInputStream(file);
} catch (FileNotFoundException e) {
e.printStackTrace();
return;
}
ZMQ.Socket router = ctx.createSocket(ZMQ.ROUTER);
router.setHWM(PIPELINE * 2);
router.bind("tcp://*:6000");
while (!Thread.currentThread().isInterrupted()) {
// First frame in each message is the sender identity
ZFrame identity = ZFrame.recvFrame(router);
if (identity.getData() == null)
break; // Shutting down, quit
// Second frame is "fetch" command
String command = router.recvStr();
assert ("fetch".equals(command));
// Third frame is chunk offset in file
int offset = Integer.parseInt(router.recvStr());
// Fourth frame is maximum chunk size
int chunkSize = Integer.parseInt(router.recvStr());
// Read chunk of data from file
byte[] data = new byte[CHUNK_SIZE];
int size;
try {
fr.skip(offset);
size = fr.read(data, 0, chunkSize);
} catch (IOException e) {
e.printStackTrace();
break;
}
// Send resulting chunk to client
ZFrame chunk = new ZFrame(Arrays.copyOf(data, size < 0 ? 0 : size));
identity.sendAndDestroy(router, ZMQ.SNDMORE);
chunk.sendAndDestroy(router);
}
try {
fr.close();
} catch (IOException e) {
e.printStackTrace();
}
}
}
}