5959 */
6060class MessageDispatcher {
6161 private static final Logger logger = Logger .getLogger (MessageDispatcher .class .getName ());
62- private static final Logger slowAckLogger = Logger .getLogger ("slow-ack" );
63- private static final Logger callbackDeliveryLogger = Logger .getLogger ("callback-delivery" );
64- private static final Logger expiryLogger = Logger .getLogger ("expiry" );
65- private static final Logger callbackExceptionsLogger = Logger .getLogger ("callback-exceptions" );
66- private static final Logger ackBatchLogger = Logger .getLogger ("ack-batch" );
67- private static final Logger subscriberFlowControlLogger =
68- Logger .getLogger ("subscriber-flow-control" );
69- private static final Logger ackNackLogger = Logger .getLogger ("ack-nack" );
62+ private LoggingUtil loggingUtil = new LoggingUtil ();
7063
7164 @ InternalApi static final double PERCENTILE_FOR_ACK_DEADLINE_UPDATES = 99.9 ;
7265 @ InternalApi static final Duration PENDING_ACKS_SEND_DELAY = Duration .ofMillis (100 );
@@ -167,15 +160,13 @@ private void forget() {
167160
168161 @ Override
169162 public void onFailure (Throwable t ) {
170- if (callbackExceptionsLogger .isLoggable (Level .WARNING )) {
171- String prefix =
172- LoggingUtil .getLogPrefix (
173- this .ackRequestData .getMessageWrapper (),
174- this .getAckRequestData ().getAckId (),
175- exactlyOnceDeliveryEnabled .get ());
176- callbackExceptionsLogger .log (
177- Level .WARNING , "pubsub:callback-exceptions - MessageReceiver exception. " + prefix , t );
178- }
163+ loggingUtil .logSubscriber (
164+ LoggingUtil .SubSytem .CALLBACK_EXCEPTIONS ,
165+ Level .WARNING ,
166+ "MessageReceiver exception." ,
167+ this .ackRequestData .getMessageWrapper (),
168+ this .ackRequestData .getAckId (),
169+ exactlyOnceDeliveryEnabled .get ());
179170 this .ackRequestData .setResponse (AckResponse .OTHER , false );
180171 pendingNacks .add (this .ackRequestData );
181172 tracer .endSubscribeProcessSpan (this .ackRequestData .getMessageWrapper (), "nack" );
@@ -186,16 +177,15 @@ public void onFailure(Throwable t) {
186177 public void onSuccess (AckReply reply ) {
187178 int ackLatency =
188179 Ints .saturatedCast ((long ) Math .ceil ((clock .millisTime () - receivedTimeMillis ) / 1000D ));
189- String logPrefix = "" ;
190- if (slowAckLogger .isLoggable (Level .FINE ) || ackNackLogger .isLoggable (Level .FINE )) {
191- logPrefix =
192- LoggingUtil .getLogPrefix (
193- this .ackRequestData .getMessageWrapper (),
194- this .ackRequestData .getAckId (),
195- exactlyOnceDeliveryEnabled .get ());
196- }
197180 if (ackLatency >= ackLatencyDistribution .getPercentile (slowAckPercentile )) {
198- slowAckLogger .log (Level .FINE , "pubsub:slow-ack - " + logPrefix );
181+ loggingUtil .logSubscriber (
182+ LoggingUtil .SubSytem .SLOW_ACK ,
183+ Level .FINE ,
184+ String .format (
185+ "Message ack duration of %d is higher than the p99 ack duration" , ackLatency ),
186+ this .ackRequestData .getMessageWrapper (),
187+ this .ackRequestData .getAckId (),
188+ exactlyOnceDeliveryEnabled .get ());
199189 }
200190
201191 switch (reply ) {
@@ -210,12 +200,24 @@ public void onSuccess(AckReply reply) {
210200 ackLatencyDistribution .record (ackLatency );
211201 tracer .endSubscribeProcessSpan (this .ackRequestData .getMessageWrapper (), "ack" );
212202 }
213- ackNackLogger .log (Level .FINE , "pubsub:ack-nack - " + logPrefix + " - Action: ACK" );
203+ loggingUtil .logSubscriber (
204+ LoggingUtil .SubSytem .ACK_NACK ,
205+ Level .FINE ,
206+ "Ack called on message." ,
207+ this .ackRequestData .getMessageWrapper (),
208+ this .ackRequestData .getAckId (),
209+ exactlyOnceDeliveryEnabled .get ());
214210 break ;
215211 case NACK :
216212 pendingNacks .add (this .ackRequestData );
217213 tracer .endSubscribeProcessSpan (this .ackRequestData .getMessageWrapper (), "nack" );
218- ackNackLogger .log (Level .FINE , "pubsub:ack-nack - " + logPrefix + " - Action: NACK" );
214+ loggingUtil .logSubscriber (
215+ LoggingUtil .SubSytem .ACK_NACK ,
216+ Level .FINE ,
217+ "Nack called on message." ,
218+ this .ackRequestData .getMessageWrapper (),
219+ this .ackRequestData .getAckId (),
220+ exactlyOnceDeliveryEnabled .get ());
219221 break ;
220222 default :
221223 throw new IllegalArgumentException (String .format ("AckReply: %s not supported" , reply ));
@@ -593,33 +595,34 @@ private void processBatch(List<OutstandingMessage> batch) {
593595 for (OutstandingMessage message : batch ) {
594596 // This is a blocking flow controller. We have already incremented messagesWaiter, so
595597 // shutdown will block on processing of all these messages anyway.
596- String logPrefix = "" ;
597- if (subscriberFlowControlLogger .isLoggable (Level .FINE )) {
598- logPrefix =
599- LoggingUtil .getLogPrefix (
600- message .messageWrapper (),
601- message .messageWrapper ().getAckId (),
602- exactlyOnceDeliveryEnabled .get ());
603- }
604598 tracer .startSubscribeConcurrencyControlSpan (message .messageWrapper ());
605599 try {
606- subscriberFlowControlLogger .log (
600+ loggingUtil .logSubscriber (
601+ LoggingUtil .SubSytem .SUBSCRIBER_FLOW_CONTROL ,
607602 Level .FINE ,
608- "pubsub:subscriber-flow-control - " + logPrefix + " - Flow controller is blocking." );
603+ "Flow controller is blocking." ,
604+ message .messageWrapper (),
605+ message .messageWrapper ().getAckId (),
606+ exactlyOnceDeliveryEnabled .get ());
609607 flowController .reserve (1 , message .messageWrapper ().getPubsubMessage ().getSerializedSize ());
610- subscriberFlowControlLogger .log (
608+ loggingUtil .logSubscriber (
609+ LoggingUtil .SubSytem .SUBSCRIBER_FLOW_CONTROL ,
611610 Level .FINE ,
612- "pubsub:subscriber-flow-control - "
613- + logPrefix
614- + " - Flow controller is done blocking." );
611+ "Flow controller is done blocking." ,
612+ message .messageWrapper (),
613+ message .messageWrapper ().getAckId (),
614+ exactlyOnceDeliveryEnabled .get ());
615615 tracer .endSubscribeConcurrencyControlSpan (message .messageWrapper ());
616616 } catch (FlowControlException unexpectedException ) {
617617 // This should be a blocking flow controller and never throw an exception.
618- subscriberFlowControlLogger .log (
618+ loggingUtil .logSubscriber (
619+ LoggingUtil .SubSytem .SUBSCRIBER_FLOW_CONTROL ,
619620 Level .FINE ,
620- "pubsub:subscriber-flow-control - "
621- + logPrefix
622- + " - Flow controller unexpected exception." );
621+ "Flow controller unexpected exception." ,
622+ message .messageWrapper (),
623+ message .messageWrapper ().getAckId (),
624+ exactlyOnceDeliveryEnabled .get (),
625+ unexpectedException );
623626 tracer .setSubscribeConcurrencyControlSpanException (
624627 message .messageWrapper (), unexpectedException );
625628 throw new IllegalStateException ("Flow control unexpected exception" , unexpectedException );
@@ -658,15 +661,6 @@ private void processOutstandingMessage(final AckHandler ackHandler) {
658661 @ Override
659662 public void run () {
660663 try {
661- String logPrefix = "" ;
662- if (expiryLogger .isLoggable (Level .FINE )
663- || callbackDeliveryLogger .isLoggable (Level .FINE )) {
664- logPrefix =
665- LoggingUtil .getLogPrefix (
666- messageWrapper ,
667- ackHandler .ackRequestData .getAckId (),
668- exactlyOnceDeliveryEnabled .get ());
669- }
670664 if (ackHandler
671665 .totalExpiration
672666 .plusSeconds (messageDeadlineSeconds .get ())
@@ -676,11 +670,23 @@ public void run() {
676670 // Don't nack it either, because we'd be nacking someone else's message.
677671 ackHandler .forget ();
678672 tracer .setSubscriberSpanExpirationResult (messageWrapper );
679- expiryLogger .log (Level .FINE , "pubsub:expiry - " + logPrefix );
673+ loggingUtil .logSubscriber (
674+ LoggingUtil .SubSytem .EXPIRY ,
675+ Level .FINE ,
676+ "Message expired." ,
677+ messageWrapper ,
678+ ackHandler .ackRequestData .getAckId (),
679+ exactlyOnceDeliveryEnabled .get ());
680680 return ;
681681 }
682682 tracer .startSubscribeProcessSpan (messageWrapper );
683- callbackDeliveryLogger .log (Level .FINE , "pubsub:callback-delivery - " + logPrefix );
683+ loggingUtil .logSubscriber (
684+ LoggingUtil .SubSytem .CALLBACK_DELIVERY ,
685+ Level .FINE ,
686+ "Message delivered." ,
687+ messageWrapper ,
688+ ackHandler .ackRequestData .getAckId (),
689+ exactlyOnceDeliveryEnabled .get ());
684690 if (shouldSetMessageFuture ()) {
685691 // This is the message future that is propagated to the user
686692 SettableApiFuture <AckResponse > messageFuture =
@@ -798,13 +804,14 @@ void processOutstandingOperations() {
798804
799805 List <AckRequestData > ackRequestDataList = new ArrayList <AckRequestData >();
800806 pendingAcks .drainTo (ackRequestDataList );
801- ackBatchLogger .log (
807+ loggingUtil .logEvent (
808+ LoggingUtil .SubSytem .ACK_BATCH ,
802809 Level .FINE ,
803- "pubsub:ack-batch - Sending {0} ACKs, {1} NACKs, {2} receipts. Exactly Once Delivery: {3}" ,
810+ "Sending {0} ACKs, {1} NACKs, {2} receipts. Exactly Once Delivery: {3}" ,
804811 new Object [] {
805812 ackRequestDataList .size (),
806813 nackRequestDataList .size (),
807- ackRequestDataList .size (),
814+ ackRequestDataReceipts .size (),
808815 exactlyOnceDeliveryEnabled .get ()
809816 });
810817
0 commit comments