@@ -103,30 +103,37 @@ public class OffsetCommitService implements Service {
103103 int heartbeatIntervalMs = config .getInt (ConsumerConfig .HEARTBEAT_INTERVAL_MS_CONFIG );
104104
105105 String clientId = config .getString (ConsumerConfig .CLIENT_ID_CONFIG );
106- LogContext logContext = new LogContext ( "[Consumer clientId=" + clientId + "] " );
106+
107107 List <String > bootstrapServers = config .getList (ConsumerConfig .BOOTSTRAP_SERVERS_CONFIG );
108108 List <InetSocketAddress > addresses =
109109 ClientUtils .parseAndValidateAddresses (bootstrapServers , ClientDnsLookup .DEFAULT );
110- ChannelBuilder channelBuilder = ClientUtils .createChannelBuilder (config , _time );
110+
111+ LogContext logContext = new LogContext ("[Consumer clientId=" + clientId + "] " );
112+
113+ ChannelBuilder channelBuilder = ClientUtils .createChannelBuilder (config , _time , logContext );
111114
112115 LOGGER .info ("Bootstrap servers config: {} | broker addresses: {}" , bootstrapServers , addresses );
113116
114117 Metadata metadata = new Metadata (retryBackoffMs , config .getLong (ConsumerConfig .METADATA_MAX_AGE_CONFIG ), logContext ,
115118 new ClusterResourceListeners ());
116119
117- metadata .bootstrap (addresses , _time . milliseconds () );
120+ metadata .bootstrap (addresses );
118121
119122 Selector selector =
120123 new Selector (config .getLong (ConsumerConfig .CONNECTIONS_MAX_IDLE_MS_CONFIG ), new Metrics (), _time ,
121124 METRIC_GRP_PREFIX , channelBuilder , logContext );
122125
123- KafkaClient kafkaClient = new NetworkClient (selector , metadata , clientId , MAX_INFLIGHT_REQUESTS_PER_CONNECTION ,
126+ KafkaClient kafkaClient = new NetworkClient (
127+ selector , metadata , clientId , MAX_INFLIGHT_REQUESTS_PER_CONNECTION ,
124128 config .getLong (ConsumerConfig .RECONNECT_BACKOFF_MS_CONFIG ),
125129 config .getLong (ConsumerConfig .RECONNECT_BACKOFF_MAX_MS_CONFIG ),
126130 config .getInt (ConsumerConfig .SEND_BUFFER_CONFIG ), config .getInt (ConsumerConfig .RECEIVE_BUFFER_CONFIG ),
127- config .getInt (ConsumerConfig .REQUEST_TIMEOUT_MS_CONFIG ), ClientDnsLookup .DEFAULT , _time , true ,
131+ config .getInt (ConsumerConfig .REQUEST_TIMEOUT_MS_CONFIG ),
132+ config .getInt (ConsumerConfig .SOCKET_CONNECTION_SETUP_TIMEOUT_MS_CONFIG ), config .getInt (ConsumerConfig .SOCKET_CONNECTION_SETUP_TIMEOUT_MAX_MS_CONFIG ),
133+ ClientDnsLookup .DEFAULT , _time , true ,
128134 new ApiVersions (), logContext );
129135
136+
130137 LOGGER .debug ("The network client active: {}" , kafkaClient .active ());
131138 LOGGER .debug ("The network client has in flight requests: {}" , kafkaClient .hasInFlightRequests ());
132139 LOGGER .debug ("The network client in flight request count: {}" , kafkaClient .inFlightRequestCount ());
0 commit comments