AE86 3 anni fa
parent
commit
f6ab99de3b

+ 1 - 0
dbsyncer-connector/src/main/java/org/dbsyncer/connector/es/ESConnector.java

@@ -187,6 +187,7 @@ public final class ESConnector extends AbstractConnector implements Connector<ES
             if (restStatus.getStatus() != RestStatus.OK.getStatus()) {
                 throw new ConnectorException(String.format("error code:%s", restStatus.getStatus()));
             }
+            result.addSuccessData(data);
         } catch (Exception e) {
             // 记录错误数据
             result.addFailData(data);

+ 1 - 0
dbsyncer-connector/src/main/java/org/dbsyncer/connector/kafka/KafkaConnector.java

@@ -83,6 +83,7 @@ public class KafkaConnector extends AbstractConnector implements Connector<Kafka
             String topic = cfg.getTopic();
             String pk = pkField.getName();
             data.forEach(row -> connectorMapper.getConnection().send(topic, String.valueOf(row.get(pk)), row));
+            result.addSuccessData(data);
         } catch (Exception e) {
             // 记录错误数据
             result.addFailData(data);