Skip to content

Commit c087124

Browse files
authored
Merge branch 'master' into docubot-edfe9a5
2 parents b3d32d7 + aaa06aa commit c087124

14 files changed

Lines changed: 22 additions & 234 deletions

File tree

‎docs/TheBook/src/main/markdown/kafkaproducer.md‎

Lines changed: 5 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -51,15 +51,13 @@ Set kafka topic name
5151
dcache.kafka.topic = billing
5252
```
5353

54-
"billing" is default. The following service level variables
55-
reference dcache.kafka.topic:
54+
"billing" is default.
5655

57-
ftp.kafka.topic = ${dcache.kafka.topic}
58-
pool.kafka.topic = ${dcache.kafka.topic}
59-
webdav.kafka.topic = ${dcache.kafka.topic}
60-
xrootd.kafka.topic = ${dcache.kafka.topic}
6156

62-
> **Note**: Kafka producer support has been removed from the NFS door. The `nfs.kafka.*` properties are obsolete and no longer have any effect.
57+
**Note** that doors no longer supports pushing messages to Kafka; the
58+
producer for doors has been removed and any `nfs/ftp/webdav.kafka.*` properties are now
59+
obsolete.
60+
6361

6462
### 3. Start Server
6563

‎modules/dcache/src/main/java/diskCacheV111/services/TransferManager.java‎

Lines changed: 0 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,6 @@
55
import diskCacheV111.util.CacheException;
66
import diskCacheV111.util.MissingResourceCacheException;
77
import diskCacheV111.util.PnfsId;
8-
import diskCacheV111.vehicles.DoorRequestInfoMessage;
98
import diskCacheV111.vehicles.DoorTransferFinishedMessage;
109
import diskCacheV111.vehicles.IpProtocolInfo;
1110
import diskCacheV111.vehicles.transferManager.CancelTransferMessage;
@@ -31,7 +30,6 @@
3130
import java.util.concurrent.ExecutorService;
3231
import java.util.concurrent.Executors;
3332
import java.util.concurrent.TimeUnit;
34-
import java.util.function.Consumer;
3533
import java.util.regex.Matcher;
3634
import java.util.regex.Pattern;
3735
import org.dcache.cells.CellStub;
@@ -40,9 +38,7 @@
4038
import org.dcache.util.CDCExecutorServiceDecorator;
4139
import org.slf4j.Logger;
4240
import org.slf4j.LoggerFactory;
43-
import org.springframework.beans.factory.annotation.Autowired;
4441
import org.springframework.beans.factory.annotation.Qualifier;
45-
import org.springframework.kafka.core.KafkaTemplate;
4642

4743

4844
/**
@@ -66,7 +62,6 @@ public abstract class TransferManager extends AbstractCellComponent
6662
private PoolManagerStub _poolManager;
6763
private CellStub _poolStub;
6864
private CellStub _billingStub;
69-
private Consumer<DoorRequestInfoMessage> _kafkaSender = (s) -> {};
7065
private boolean _overwrite;
7166
private int _maxNumberOfDeleteRetries;
7267
// this is the timer which will timeout the
@@ -371,10 +366,6 @@ public CellStub getBillingStub() {
371366
return _billingStub;
372367
}
373368

374-
public Consumer<DoorRequestInfoMessage> getKafkaSender() {
375-
return _kafkaSender;
376-
}
377-
378369
public String getIoQueueName() {
379370
return _ioQueueName;
380371
}
@@ -396,11 +387,6 @@ public void setBilling(CellStub billingStub) {
396387
_billingStub = billingStub;
397388
}
398389

399-
@Autowired(required = false)
400-
public void setTransferTemplate(KafkaTemplate kafkaTemplate) {
401-
_kafkaSender = kafkaTemplate::sendDefault;
402-
}
403-
404390
public void setPoolManager(PoolManagerStub poolManager) {
405391
_poolManager = poolManager;
406392
}

‎modules/dcache/src/main/java/diskCacheV111/services/TransferManagerHandler.java‎

Lines changed: 0 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,6 @@
66
import static org.dcache.namespace.FileAttribute.STORAGECLASS;
77
import static org.dcache.namespace.FileAttribute.STORAGEINFO;
88

9-
import com.google.common.base.Throwables;
109
import com.google.common.collect.ImmutableMap;
1110
import com.google.common.collect.Sets;
1211
import diskCacheV111.util.CacheException;
@@ -61,7 +60,6 @@
6160
import org.dcache.vehicles.PnfsGetFileAttributes;
6261
import org.slf4j.Logger;
6362
import org.slf4j.LoggerFactory;
64-
import org.springframework.kafka.KafkaException;
6563

6664
public class TransferManagerHandler extends AbstractMessageCallback<Message> {
6765

@@ -742,13 +740,6 @@ void sendDoorRequestInfo(int code, String msg) {
742740
info.setResult(code, msg);
743741
LOGGER.debug("Sending info: {}", info);
744742
manager.getBillingStub().notify(info);
745-
746-
try {
747-
manager.getKafkaSender().accept(info);
748-
} catch (KafkaException | org.apache.kafka.common.KafkaException e) {
749-
LOGGER.warn("Failed to send message to kafka: {} ",
750-
Throwables.getRootCause(e).getMessage());
751-
}
752743
}
753744

754745
public void timeout() {

‎modules/dcache/src/main/java/org/dcache/pool/classic/DefaultPostTransferService.java‎

Lines changed: 0 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,6 @@
2020
import static com.google.common.net.InetAddresses.toUriString;
2121
import static org.dcache.util.Exceptions.messageOrClassName;
2222

23-
import com.google.common.base.Throwables;
2423
import com.google.common.util.concurrent.MoreExecutors;
2524
import com.google.common.util.concurrent.ThreadFactoryBuilder;
2625
import diskCacheV111.util.CacheException;
@@ -40,7 +39,6 @@
4039
import java.util.concurrent.ExecutorService;
4140
import java.util.concurrent.Executors;
4241
import java.util.concurrent.TimeUnit;
43-
import java.util.function.Consumer;
4442
import org.dcache.cells.CellStub;
4543
import org.dcache.pool.movers.Mover;
4644
import org.dcache.pool.movers.TransferLifeCycle;
@@ -56,11 +54,7 @@
5654
import org.dcache.vehicles.FileAttributes;
5755
import org.slf4j.Logger;
5856
import org.slf4j.LoggerFactory;
59-
import org.springframework.beans.factory.annotation.Autowired;
60-
import org.springframework.beans.factory.annotation.Qualifier;
6157
import org.springframework.beans.factory.annotation.Required;
62-
import org.springframework.kafka.KafkaException;
63-
import org.springframework.kafka.core.KafkaTemplate;
6458

6559
public class DefaultPostTransferService extends AbstractCellComponent implements
6660
PostTransferService, CellInfoProvider {
@@ -77,9 +71,6 @@ public class DefaultPostTransferService extends AbstractCellComponent implements
7771
private ChecksumModule _checksumModule;
7872
private CellStub _door;
7973

80-
private Consumer<MoverInfoMessage> _kafkaSender = (s) -> {
81-
};
82-
8374
private TransferLifeCycle transferLifeCycle;
8475

8576
@Required
@@ -97,12 +88,6 @@ public void setChecksumModule(ChecksumModule checksumModule) {
9788
_checksumModule = checksumModule;
9889
}
9990

100-
@Autowired(required = false)
101-
@Qualifier("transfer")
102-
public void setKafkaTemplate(KafkaTemplate kafkaTemplate) {
103-
_kafkaSender = kafkaTemplate::sendDefault;
104-
}
105-
10691
public void setTransferLifeCycle(TransferLifeCycle transferLifeCycle) {
10792
this.transferLifeCycle = transferLifeCycle;
10893
}
@@ -169,12 +154,6 @@ public void execute(final Mover<?> mover,
169154

170155
private void sendBillingInfo(MoverInfoMessage moverInfoMessage) {
171156
_billing.notify(moverInfoMessage);
172-
173-
try {
174-
_kafkaSender.accept(moverInfoMessage);
175-
} catch (KafkaException | org.apache.kafka.common.KafkaException e) {
176-
LOGGER.warn("Failed to send message to kafka: {} ", Throwables.getRootCause(e).getMessage());
177-
}
178157
}
179158

180159
public MoverInfoMessage generateBillingMessage(Mover<?> mover, long fileSize) {

‎modules/dcache/src/main/java/org/dcache/pool/classic/PoolV4.java‎

Lines changed: 0 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,6 @@
77
import static java.util.stream.Collectors.toList;
88

99
import com.google.common.base.Splitter;
10-
import com.google.common.base.Throwables;
1110
import com.google.common.collect.ImmutableMap;
1211
import com.google.common.collect.ImmutableSet;
1312
import com.google.common.net.InetAddresses;
@@ -87,7 +86,6 @@
8786
import java.util.concurrent.Callable;
8887
import java.util.concurrent.Executor;
8988
import java.util.concurrent.ThreadFactory;
90-
import java.util.function.Consumer;
9189
import java.util.stream.Stream;
9290
import org.dcache.alarms.AlarmMarkerFactory;
9391
import org.dcache.alarms.PredefinedAlarm;
@@ -123,11 +121,7 @@
123121
import org.dcache.vehicles.FileAttributes;
124122
import org.slf4j.Logger;
125123
import org.slf4j.LoggerFactory;
126-
import org.springframework.beans.factory.annotation.Autowired;
127-
import org.springframework.beans.factory.annotation.Qualifier;
128124
import org.springframework.beans.factory.annotation.Required;
129-
import org.springframework.kafka.KafkaException;
130-
import org.springframework.kafka.core.KafkaTemplate;
131125

132126
public class PoolV4
133127
extends AbstractCellComponent
@@ -209,9 +203,6 @@ public class PoolV4
209203

210204
private Executor _executor;
211205

212-
private Consumer<RemoveFileInfoMessage> _kafkaSender = (s) -> {
213-
};
214-
215206
private ThreadFactory _threadFactory;
216207

217208
// Hot file monitoring
@@ -239,12 +230,6 @@ protected void assertNotRunning(String error) {
239230
checkState(!_running, error);
240231
}
241232

242-
@Autowired(required = false)
243-
@Qualifier("remove")
244-
public void setKafkaTemplate(KafkaTemplate kafkaTemplate) {
245-
_kafkaSender = kafkaTemplate::sendDefault;
246-
}
247-
248233
@Required
249234
public void setPoolName(String name) {
250235
assertNotRunning("Cannot change pool name after initialisation");
@@ -596,13 +581,6 @@ public void stateChanged(StateChangeEvent event) {
596581
msg.setSubject(Subjects.ROOT);
597582
msg.setResult(0, event.getWhy());
598583
_billingStub.notify(msg);
599-
600-
try {
601-
_kafkaSender.accept(msg);
602-
} catch (KafkaException | org.apache.kafka.common.KafkaException e) {
603-
LOGGER.warn("Failed to send message to kafka: {} ",
604-
Throwables.getRootCause(e).getMessage());
605-
}
606584
}
607585
}
608586
}

‎modules/dcache/src/main/java/org/dcache/pool/nearline/NearlineStorageHandler.java‎

Lines changed: 0 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,6 @@
2727
import static org.dcache.util.Exceptions.messageOrClassName;
2828

2929
import com.google.common.base.Functions;
30-
import com.google.common.base.Throwables;
3130
import com.google.common.util.concurrent.AsyncFunction;
3231
import com.google.common.util.concurrent.Futures;
3332
import com.google.common.util.concurrent.ListenableFuture;
@@ -80,7 +79,6 @@
8079
import java.util.concurrent.ScheduledFuture;
8180
import java.util.concurrent.TimeUnit;
8281
import java.util.concurrent.atomic.AtomicReference;
83-
import java.util.function.Consumer;
8482
import java.util.function.Predicate;
8583
import java.util.stream.Collectors;
8684
import javax.annotation.Nonnull;
@@ -119,11 +117,7 @@
119117
import org.dcache.vehicles.FileAttributes;
120118
import org.slf4j.Logger;
121119
import org.slf4j.LoggerFactory;
122-
import org.springframework.beans.factory.annotation.Autowired;
123-
import org.springframework.beans.factory.annotation.Qualifier;
124120
import org.springframework.beans.factory.annotation.Required;
125-
import org.springframework.kafka.KafkaException;
126-
import org.springframework.kafka.core.KafkaTemplate;
127121

128122
/**
129123
* Entry point to and management interface for the nearline storage subsystem.
@@ -173,10 +167,6 @@ public record QueueStat(int queued, int active, int canceled) {
173167

174168
private CellAddressCore cellAddress;
175169

176-
private Consumer<StorageInfoMessage> _kafkaSender = (s) -> {
177-
};
178-
179-
180170
@Override
181171
public void setCellAddress(CellAddressCore address) {
182172
cellAddress = address;
@@ -188,12 +178,6 @@ public void setScheduledExecutor(ScheduledExecutorService executor) {
188178
}
189179

190180

191-
@Autowired(required = false)
192-
@Qualifier("hsm")
193-
public void setKafkaTemplate(KafkaTemplate kafkaTemplate) {
194-
_kafkaSender = kafkaTemplate::sendDefault;
195-
}
196-
197181
@Required
198182
public void setExecutor(ListeningExecutorService executor) {
199183
this.executor = requireNonNull(executor);
@@ -1187,11 +1171,6 @@ private void done(@Nullable Throwable cause) {
11871171
addFromNearlineStorage(infoMsg, storage);
11881172

11891173
billingStub.notify(infoMsg);
1190-
try {
1191-
_kafkaSender.accept(infoMsg);
1192-
} catch (KafkaException | org.apache.kafka.common.KafkaException e) {
1193-
LOGGER.warn("Failed to send message to kafka: {} ", Throwables.getRootCause(e).getMessage());
1194-
}
11951174
flushRequests.removeAndCallback(pnfsId, cause);
11961175
}
11971176

@@ -1421,11 +1400,6 @@ private void done(@Nullable Throwable cause) {
14211400
addFromNearlineStorage(infoMsg, storage);
14221401

14231402
billingStub.notify(infoMsg);
1424-
try {
1425-
_kafkaSender.accept(infoMsg);
1426-
} catch (KafkaException | org.apache.kafka.common.KafkaException e) {
1427-
LOGGER.warn("Failed to send message to kafka: {} ", Throwables.getRootCause(e).getMessage());
1428-
}
14291403
stageRequests.removeAndCallback(pnfsId, cause);
14301404
}
14311405

‎modules/dcache/src/main/java/org/dcache/util/jetty/RateLimitedHandlerList.java‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -240,6 +240,10 @@ public void handle(String target, Request baseRequest, HttpServletRequest reques
240240
return;
241241
}
242242

243+
244+
245+
// REVISIT: do we want to keep per-client rate limiter or only for authentication errors?
246+
/*
243247
if (!getClientRateLimiter(client).tryAcquire()) {
244248
LOGGER.debug("Blocking client with too many requests {}", client);
245249
response.setStatus(HttpStatus.TOO_MANY_REQUESTS_429);
@@ -257,7 +261,7 @@ public void handle(String target, Request baseRequest, HttpServletRequest reques
257261
baseRequest.setHandled(true);
258262
return;
259263
}
260-
264+
*/
261265
Handler[] handlers = this.getHandlers();
262266
if (handlers != null && this.isStarted()) {
263267
for (Handler handler : handlers) {

‎modules/dcache/src/main/resources/diskCacheV111/services/transfer-manager.xml‎

Lines changed: 0 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -82,33 +82,4 @@
8282
<property name="overwrite" value="false"/>
8383
</bean>
8484

85-
<beans profile="kafka-true">
86-
87-
<bean id="listener" class="org.dcache.kafka.LoggingProducerListener"/>
88-
89-
<bean id="kafka-configs"
90-
class="org.dcache.util.configuration.ConfigurationMapFactoryBean">
91-
<property name="prefix" value="transfermanagers.kafka.producer.configs"/>
92-
<property name="staticEnvironment">
93-
<map>
94-
<entry key="bootstrap.servers" value="${transfermanagers.kafka.producer.bootstrap.servers}"/>
95-
<entry key="key.serializer" value="org.apache.kafka.common.serialization.StringSerializer"/>
96-
<entry key="value.serializer" value="org.dcache.notification.DoorRequestMessageSerializer"/>
97-
<entry key="client.id" value="${transfermanagers.cell.name}@${dcache.domain.name}"/>
98-
</map>
99-
</property>
100-
</bean>
101-
102-
<bean id="billing-template" class="org.springframework.kafka.core.KafkaTemplate">
103-
<constructor-arg>
104-
<bean class="org.springframework.kafka.core.DefaultKafkaProducerFactory">
105-
<constructor-arg name="configs" ref="kafka-configs"/>
106-
</bean>
107-
</constructor-arg>
108-
<property name="defaultTopic" value="${transfermanagers.kafka.topic}"/>
109-
<property name="producerListener" ref="listener"/>
110-
</bean>
111-
112-
</beans>
113-
11485
</beans>

0 commit comments

Comments
 (0)