|
37 | 37 | import io.questdb.client.std.Decimal128; |
38 | 38 | import io.questdb.client.std.Decimal256; |
39 | 39 | import io.questdb.client.std.Decimal64; |
| 40 | +import io.questdb.client.std.Misc; |
40 | 41 | import io.questdb.client.std.ObjList; |
41 | 42 | import io.questdb.client.std.Unsafe; |
42 | 43 | import io.questdb.client.std.bytes.DirectByteSlice; |
@@ -69,10 +70,10 @@ public class QwpUdpSender implements Sender { |
69 | 70 | private static final int VARINT_INT_UPPER_BOUND = 5; |
70 | 71 | private final UdpLineChannel channel; |
71 | 72 | private final QwpColumnWriter columnWriter = new QwpColumnWriter(); |
72 | | - private final NativeSegmentList datagramSegments = new NativeSegmentList(); |
73 | | - private final NativeBufferWriter headerBuffer = new NativeBufferWriter(); |
| 73 | + private final NativeSegmentList datagramSegments; |
| 74 | + private final NativeBufferWriter headerBuffer; |
74 | 75 | private final int maxDatagramSize; |
75 | | - private final SegmentedNativeBufferWriter payloadWriter = new SegmentedNativeBufferWriter(); |
| 76 | + private final SegmentedNativeBufferWriter payloadWriter; |
76 | 77 | private final CharSequenceObjHashMap<QwpTableBuffer> tableBuffers; |
77 | 78 | private final CharSequenceObjHashMap<TableHeadroomState> tableHeadroomStates; |
78 | 79 | private final boolean trackDatagramEstimate; |
@@ -105,9 +106,28 @@ public QwpUdpSender(NetworkFacade nf, int interfaceIPv4, int sendToAddress, int |
105 | 106 | } |
106 | 107 |
|
107 | 108 | public QwpUdpSender(NetworkFacade nf, int interfaceIPv4, int sendToAddress, int port, int ttl, int maxDatagramSize) { |
108 | | - this.channel = new UdpLineChannel(nf, interfaceIPv4, sendToAddress, port, ttl); |
109 | | - this.tableHeadroomStates = new CharSequenceObjHashMap<>(); |
110 | | - this.tableBuffers = new CharSequenceObjHashMap<>(); |
| 109 | + NativeSegmentList segments = null; |
| 110 | + NativeBufferWriter header = null; |
| 111 | + SegmentedNativeBufferWriter payload = null; |
| 112 | + UdpLineChannel ch = null; |
| 113 | + try { |
| 114 | + segments = new NativeSegmentList(); |
| 115 | + header = new NativeBufferWriter(); |
| 116 | + payload = new SegmentedNativeBufferWriter(); |
| 117 | + ch = new UdpLineChannel(nf, interfaceIPv4, sendToAddress, port, ttl); |
| 118 | + this.tableHeadroomStates = new CharSequenceObjHashMap<>(); |
| 119 | + this.tableBuffers = new CharSequenceObjHashMap<>(); |
| 120 | + } catch (Throwable t) { |
| 121 | + Misc.free(ch); |
| 122 | + Misc.free(payload); |
| 123 | + Misc.free(header); |
| 124 | + Misc.free(segments); |
| 125 | + throw t; |
| 126 | + } |
| 127 | + this.channel = ch; |
| 128 | + this.datagramSegments = segments; |
| 129 | + this.headerBuffer = header; |
| 130 | + this.payloadWriter = payload; |
111 | 131 | this.maxDatagramSize = maxDatagramSize; |
112 | 132 | this.trackDatagramEstimate = maxDatagramSize > 0; |
113 | 133 | } |
|
0 commit comments