Movatterモバイル変換


[0]ホーム

URL:


Skip to content

Navigation Menu

Sign in
Appearance settings

Search code, repositories, users, issues, pull requests...

Provide feedback

We read every piece of feedback, and take your input very seriously.

Saved searches

Use saved searches to filter your results more quickly

Sign up
Appearance settings

Commit829189b

Browse files
committed
updates
1 parent7ded463 commit829189b

File tree

1 file changed

+15
-24
lines changed

1 file changed

+15
-24
lines changed

‎src/main/java/com/owl/kafka/proxy/server/biz/pull/PullCenter.java‎

Lines changed: 15 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -41,8 +41,6 @@ public class PullCenter{
4141

4242
privatefinalPullRequestHoldServicepullRequestHoldService =newPullRequestHoldService();
4343

44-
privatefinalByteBufAllocatorallocator =newPooledByteBufAllocator(true);
45-
4644
publicvoidputMessage(ConsumerRecord<byte[],byte[]>record)throwsInterruptedException{
4745
this.pullQueue.put(record);
4846
this.pullRequestHoldService.notifyMessageArriving();
@@ -73,14 +71,10 @@ private boolean poll(Packet packet) {
7371
Packetone =retryQueue.peek();
7472
if(one !=null){
7573
retryQueue.poll();
76-
ByteBufbuffer =allocator.directBuffer(one.getBody().length +packet.getBody().length);
77-
try {
78-
buffer.writeBytes(packet.getBody());
79-
buffer.writeBytes(one.getBody());
80-
packet.setBody(buffer.array());
81-
}finally {
82-
buffer.release();
83-
}
74+
ByteBufferbuffer =ByteBuffer.allocate(one.getBody().length +packet.getBody().length);
75+
buffer.put(packet.getBody());
76+
buffer.put(one.getBody());
77+
packet.setBody(buffer.array());
8478
polled =true;
8579
}else{
8680
ConsumerRecord<byte[],byte[]>record =pullQueue.poll();
@@ -90,20 +84,17 @@ private boolean poll(Packet packet) {
9084
byte[]headerInBytes =SerializerImpl.getFastJsonSerializer().serialize(header);
9185
//
9286
intcapacity =4 +headerInBytes.length +4 +record.key().length +4 +record.value().length;
93-
ByteBufbuffer =allocator.directBuffer(capacity +packet.getBody().length);
94-
try {
95-
buffer.writeBytes(packet.getBody());
96-
buffer.writeInt(headerInBytes.length);
97-
buffer.writeBytes(headerInBytes);
98-
buffer.writeInt(record.key().length);
99-
buffer.writeBytes(record.key());
100-
buffer.writeInt(record.value().length);
101-
buffer.writeBytes(record.value());
102-
//
103-
packet.setBody(buffer.array());
104-
}finally {
105-
buffer.release();
106-
}
87+
ByteBufferbuffer =ByteBuffer.allocate(capacity +packet.getBody().length);
88+
89+
buffer.put(packet.getBody());
90+
buffer.putInt(headerInBytes.length);
91+
buffer.put(headerInBytes);
92+
buffer.putInt(record.key().length);
93+
buffer.put(record.key());
94+
buffer.putInt(record.value().length);
95+
buffer.put(record.value());
96+
//
97+
packet.setBody(buffer.array());
10798
polled =true;
10899
}
109100
}

0 commit comments

Comments
 (0)

[8]ページ先頭

©2009-2025 Movatter.jp