fuweng11 commented on code in PR #7664:
URL: https://github.com/apache/inlong/pull/7664#discussion_r1143159740


##########
inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/StreamSinkServiceImpl.java:
##########
@@ -695,6 +704,57 @@ public List<SinkField> parseFields(String fieldsJson) {
         }
     }
 
+    @Override
+    public void addFieldForSink(AddFieldsRequest fieldsRequest, String 
sourceType, InlongGroupEntity groupEntity,
+            InlongStreamEntity streamEntity) {
+        AtomicBoolean isNeedAddField = new AtomicBoolean(false);
+        String groupId = groupEntity.getInlongGroupId();
+        String streamId = streamEntity.getInlongStreamId();
+        String defaultOperator = 
groupEntity.getInCharges().split(InlongConstants.COMMA)[0];
+        // add fields for StreamSinkField
+        List<StreamSinkEntity> sinkEntityList = 
sinkMapper.selectByRelatedId(groupId, streamId);
+        sinkEntityList.forEach(sink -> {
+            String sinkType = sink.getSinkType();
+            List<SinkField> toAddFields = 
fieldsRequest.getFields().stream().map(streamField -> {
+                SinkField sinkField = new SinkField();
+                sinkField.setSinkType(sink.getSinkType());
+                sinkField.setFieldName(streamField.getFieldName());
+                sinkField.setFieldType(fieldTypeUtils.getSinkField(sinkType, 
streamField.getFieldType()));
+                sinkField.setFieldComment(streamField.getFieldComment());
+                sinkField.setSourceFieldName(streamField.getFieldName());
+                
sinkField.setSourceFieldType(fieldTypeUtils.getStreamField(sourceType, 
streamField.getFieldType()));
+                return sinkField;
+            }).collect(Collectors.toList());
+            List<StreamSinkFieldEntity> existsFieldList = 
sinkFieldMapper.selectBySinkId(sink.getId());
+            List<SinkField> sinkFields = new ArrayList<>();
+            if (CollectionUtils.isNotEmpty(existsFieldList)) {
+                sinkFields = 
CommonBeanUtils.copyListProperties(existsFieldList, SinkField::new);
+            }
+            Set<String> existsNames = sinkFields.stream()
+                    .map(field -> 
field.getFieldName().toLowerCase(Locale.ROOT))
+                    .collect(Collectors.toSet());
+            for (SinkField fieldInfo : toAddFields) {
+                String tobeAddFieldName = 
fieldInfo.getFieldName().toLowerCase(Locale.ROOT);
+                if (existsNames.contains(tobeAddFieldName)) {
+                    LOGGER.error("sink field {} already exist for sinkId {}", 
fieldInfo.getFieldName(), sink.getId());

Review Comment:
   No, because the agent may report multiple times



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to