Updated to rephrase question with additional information
We have two annotations:
CustomLoggingPollableStreamListener
Both are implemented using aspects with Spring AOP.
CustomLogging annotation:
@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
public @interface CustomLogging {
}
CustomLoggingAspect class:
@Aspect
@Component
@Slf4j
@Order(value = 1)
public class CustomLoggingAspect {
@Before("@annotation(customLogging)")
public void addCustomLogging(CustomLogging customLogging) {
log.info("Logging some information");
}
}
PollableStreamListener annotation:
@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
public @interface PollableStreamListener {
}
PollableStreamListenerAspect class:
@Aspect
@Component
@Slf4j
public class PollableStreamListenerAspect {
private final ExecutorService executor = Executors.newFixedThreadPool(1);
private volatile boolean paused = false;
@Around(value = "@annotation(pollableStreamListener) && args(dataCapsule,..)")
public void receiveMessage(ProceedingJoinPoint joinPoint,
PollableStreamListener pollableStreamListener, Object dataCapsule) {
if (dataCapsule instanceof Message) {
Message<?> message = (Message<?>) dataCapsule;
AcknowledgmentCallback callback = StaticMessageHeaderAccessor
.getAcknowledgmentCallback(message);
callback.noAutoAck();
if (!paused) {
// The separate thread is not busy with a previous message, so process this message:
Runnable runnable = () -> {
try {
paused = true;
// Call method to process this Kafka message
joinPoint.proceed();
callback.acknowledge(Status.ACCEPT);
} catch (Throwable e) {
callback.acknowledge(Status.REJECT);
throw new PollableStreamListenerException(e);
} finally {
paused = false;
}
};
executor.submit(runnable);
} else {
// The separate thread is busy with a previous message, so re-queue this message for later:
callback.acknowledge(Status.REQUEUE);
log.info("Re-queue");
}
}
}
}
We have a class called CleanupController which regularly executes according to a schedule.
CleanupController class:
@Scheduled(fixedDelayString = "${app.pollable-consumer.time-interval}")
public void pollForDeletionRequest() {
log.trace("Polling for new messages");
cleanupInput.poll(cleanupSubmissionService::submitDeletion);
}
When the schedule executes, it calls a method in another class that is annotated with both PollableStreamListener and CustomLogging. I've added a Thread.sleep() to imitate the method taking a while to execute.
@PollableStreamListener
@CustomLogging
public void submitDeletion(Message<?> received) {
try {
log.info("Starting processing");
Thread.sleep(10000);
log.info("Finished processing");
} catch (Exception e) {
log.info("Error", e);
}
}
The problem I'm facing is that the output produced by CustomLogging is printed each time we poll for new messages using the @Schedule, but I only want it to print if the annotated method is actually executed (which may happen now, or may happen in the future depending on whether another message is currently being processed). This leads to confusing log messages because it implies that the message is being processed now when in fact it has been re-queued for future execution.
Is there some way that we can make these annotations work well together so that the CustomLogging output only happens if the annotated method executes?
Update to use @Order on PollableStreamListener
As per the suggestion of @dunni, I made the following changes to the original example above.
Set the order of 1 on the PollableStreamListenerAspect:
@Aspect
@Component
@Slf4j
@Order(value = 1)
public class PollableStreamListenerAspect {
...
}
Increase the order to 2 for CustomLoggingAspect:
@Aspect
@Component
@Slf4j
@Order(value = 2)
public class CustomLoggingAspect {
...
}
I found that after making these changes that the polling fails to detect new requests at all. It is the change on the PollableStreamListenerAspect that introduced this issue (I commented out that line and re-ran it, and things behaved as they did before).
Update to use @Order(value = Ordered.HIGHEST_PRECEDENCE) on PollableStreamListener
I've updated the PollableStreamListener to use the HIGHEST_PRECEDENCE and update the @Around value:
@Aspect
@Component
@Slf4j
@Order(value = Ordered.HIGHEST_PRECEDENCE)
public class PollableStreamListenerAspect {
private final ExecutorService executor = Executors.newFixedThreadPool(1);
private volatile boolean paused = false;
@Around(value = "@annotation(au.com.brolly.commons.stream.PollableStreamListener)")
public void receiveMessage(ProceedingJoinPoint joinPoint) {
if (!paused) {
// The separate thread is not busy with a previous message, so process this message:
Runnable runnable = () -> {
try {
paused = true;
// Call method to process this Kafka message
joinPoint.proceed();
} catch (Throwable e) {
e.printStackTrace();
throw new PollableStreamListenerException(e);
} finally {
paused = false;
}
};
executor.submit(runnable);
} else {
// The separate thread is busy with a previous message, so re-queue this message for later:
log.info("Re-queue");
}
}
}
This is partially working now. When I send a Kafka message, it gets processed, and the logging from the CustomLogging annotation only prints if another Kafka message isn't being processed. So far so good.
The next challenge is to get the @Around to accept the Message that is being provided via Kafka. I attempted this using the above example with the following lines changed:
@Around(value = "@annotation(au.com.brolly.commons.stream.PollableStreamListener) && args(dataCapsule,..)")
public void receiveMessage(ProceedingJoinPoint joinPoint, Object dataCapsule) {
...
}
The server starts correctly, but when I publish a Kafka message then I get the following exception:
2021-04-22 10:38:00,055 ERROR [scheduling-1] org.springframework.core.log.LogAccessor: org.springframework.messaging.MessageHandlingException: nested exception is java.lang.IllegalStateException: Required to bind 2 arguments, but only bound 1 (JoinPointMatch was NOT bound in invocation), failedMessage=GenericMessage...
at org.springframework.cloud.stream.binder.DefaultPollableMessageSource.doHandleMessage(DefaultPollableMessageSource.java:330)
at org.springframework.cloud.stream.binder.DefaultPollableMessageSource.handle(DefaultPollableMessageSource.java:361)
at org.springframework.cloud.stream.binder.DefaultPollableMessageSource.poll(DefaultPollableMessageSource.java:219)
at org.springframework.cloud.stream.binder.DefaultPollableMessageSource.poll(DefaultPollableMessageSource.java:200)
at org.springframework.cloud.stream.binder.DefaultPollableMessageSource.poll(DefaultPollableMessageSource.java:68)
at xyx.pollForDeletionRequest(CleanupController.java:35)
at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
at java.base/java.lang.reflect.Method.invoke(Method.java:566)
at org.springframework.scheduling.support.ScheduledMethodRunnable.run(ScheduledMethodRunnable.java:84)
at org.springframework.scheduling.support.DelegatingErrorHandlingRunnable.run(DelegatingErrorHandlingRunnable.java:54)
at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
at java.base/java.util.concurrent.FutureTask.runAndReset(FutureTask.java:305)
at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:305)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
at java.base/java.lang.Thread.run(Thread.java:834)
Caused by: java.lang.IllegalStateException: Required to bind 2 arguments, but only bound 1 (JoinPointMatch was NOT bound in invocation)
at org.springframework.aop.aspectj.AbstractAspectJAdvice.argBinding(AbstractAspectJAdvice.java:596)
at org.springframework.aop.aspectj.AbstractAspectJAdvice.invokeAdviceMethod(AbstractAspectJAdvice.java:624)
at org.springframework.aop.aspectj.AspectJAroundAdvice.invoke(AspectJAroundAdvice.java:72)
at org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:175)
at org.springframework.aop.framework.CglibAopProxy$CglibMethodInvocation.proceed(CglibAopProxy.java:750)
at org.springframework.aop.framework.CglibAopProxy$DynamicAdvisedInterceptor.intercept(CglibAopProxy.java:692)
at xyz.CleanupSubmissionServiceImpl$$EnhancerBySpringCGLIB$$8737f6f8.submitDeletion(<generated>)
at org.springframework.cloud.stream.binder.DefaultPollableMessageSource.doHandleMessage(DefaultPollableMessageSource.java:327)
... 17 more