This is an automated email from the ASF dual-hosted git repository.
tenthe pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to refs/heads/dev by this push:
new 2d44c44cc8 feat(#4721): add CSV import progress tracking (#4722)
2d44c44cc8 is described below
commit 2d44c44cc896f63a1008161369e4af286ebba2e1
Author: Philipp Zehnder <[email protected]>
AuthorDate: Wed Jul 15 13:09:49 2026 +0200
feat(#4721): add CSV import progress tracking (#4722)
---
.../datalake/importer/CsvImportJobStartResult.java | 39 +++++
.../model/datalake/importer/CsvImportJobState.java | 25 ++++
.../datalake/importer/CsvImportJobStatus.java | 89 +++++++++++
.../datalake/importer/CsvImportPreviewResult.java | 9 ++
.../rest/impl/datalake/DataLakeResource.java | 10 +-
.../importer/CsvDataLakeImportService.java | 71 ++++++++-
.../datalake/importer/CsvImportJobManager.java | 162 +++++++++++++++++++++
.../impl/datalake/importer/CsvImportParser.java | 29 +++-
.../datalake/importer/CsvImportUploadStorage.java | 20 ++-
.../datalake/importer/DataLakeImportResource.java | 21 ++-
.../rest/impl/datalake/DataLakeResourceTest.java | 17 +++
.../importer/CsvDataLakeImportServiceTest.java | 65 +++++++++
.../datalake/importer/CsvImportParserTest.java | 1 +
ui/STYLEGUIDE.md | 14 ++
.../support/utils/dataset/DataLakeSeedUtils.ts | 69 ++++++++-
ui/deployment/i18n/de.json | 3 +-
ui/deployment/i18n/en.json | 3 +-
ui/deployment/i18n/pl.json | 3 +-
.../src/lib/apis/datalake-rest.service.ts | 18 ++-
.../src/lib/model/datalake/csv-import.model.ts | 16 ++
.../progress-bar/progress-bar.component.html | 37 +++++
.../progress-bar/progress-bar.component.scss} | 37 ++---
.../progress-bar/progress-bar.component.ts | 57 ++++++++
.../streampipes/shared-ui/src/public-api.ts | 1 +
.../csv-import-dialog.component.html | 2 +
.../csv-import-dialog.component.ts | 92 +++++++++++-
.../csv-import-upload-state.component.html | 19 ++-
.../csv-import-upload-state.component.scss | 19 ++-
.../csv-import-upload-state.component.ts | 10 +-
29 files changed, 874 insertions(+), 84 deletions(-)
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/importer/CsvImportJobStartResult.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/importer/CsvImportJobStartResult.java
new file mode 100644
index 0000000000..0b066ba223
--- /dev/null
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/importer/CsvImportJobStartResult.java
@@ -0,0 +1,39 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.model.datalake.importer;
+
+public class CsvImportJobStartResult {
+
+ private String jobId;
+
+ public CsvImportJobStartResult() {
+ }
+
+ public CsvImportJobStartResult(String jobId) {
+ this.jobId = jobId;
+ }
+
+ public String getJobId() {
+ return jobId;
+ }
+
+ public void setJobId(String jobId) {
+ this.jobId = jobId;
+ }
+}
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/importer/CsvImportJobState.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/importer/CsvImportJobState.java
new file mode 100644
index 0000000000..b97de417e0
--- /dev/null
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/importer/CsvImportJobState.java
@@ -0,0 +1,25 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.model.datalake.importer;
+
+public enum CsvImportJobState {
+ RUNNING,
+ SUCCEEDED,
+ FAILED
+}
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/importer/CsvImportJobStatus.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/importer/CsvImportJobStatus.java
new file mode 100644
index 0000000000..495930ab65
--- /dev/null
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/importer/CsvImportJobStatus.java
@@ -0,0 +1,89 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.model.datalake.importer;
+
+import java.util.ArrayList;
+import java.util.List;
+
+public class CsvImportJobStatus {
+
+ private String jobId;
+ private CsvImportJobState state;
+ private int processedRows;
+ private int totalRows;
+ private int progress;
+ private CsvImportResult result;
+ private List<CsvImportValidationMessage> validationMessages = new
ArrayList<>();
+
+ public String getJobId() {
+ return jobId;
+ }
+
+ public void setJobId(String jobId) {
+ this.jobId = jobId;
+ }
+
+ public CsvImportJobState getState() {
+ return state;
+ }
+
+ public void setState(CsvImportJobState state) {
+ this.state = state;
+ }
+
+ public int getProcessedRows() {
+ return processedRows;
+ }
+
+ public void setProcessedRows(int processedRows) {
+ this.processedRows = processedRows;
+ }
+
+ public int getTotalRows() {
+ return totalRows;
+ }
+
+ public void setTotalRows(int totalRows) {
+ this.totalRows = totalRows;
+ }
+
+ public int getProgress() {
+ return progress;
+ }
+
+ public void setProgress(int progress) {
+ this.progress = progress;
+ }
+
+ public CsvImportResult getResult() {
+ return result;
+ }
+
+ public void setResult(CsvImportResult result) {
+ this.result = result;
+ }
+
+ public List<CsvImportValidationMessage> getValidationMessages() {
+ return validationMessages;
+ }
+
+ public void setValidationMessages(List<CsvImportValidationMessage>
validationMessages) {
+ this.validationMessages = validationMessages;
+ }
+}
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/importer/CsvImportPreviewResult.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/importer/CsvImportPreviewResult.java
index 4bb9b84cf8..f2489410a1 100644
---
a/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/importer/CsvImportPreviewResult.java
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/datalake/importer/CsvImportPreviewResult.java
@@ -31,6 +31,7 @@ public class CsvImportPreviewResult {
private List<CsvImportColumn> columns = new ArrayList<>();
private EventSchema guessedEventSchema;
private List<String> timestampCandidates = new ArrayList<>();
+ private int totalRows;
private boolean valid;
private List<CsvImportValidationMessage> validationMessages = new
ArrayList<>();
@@ -82,6 +83,14 @@ public class CsvImportPreviewResult {
this.timestampCandidates = timestampCandidates;
}
+ public int getTotalRows() {
+ return totalRows;
+ }
+
+ public void setTotalRows(int totalRows) {
+ this.totalRows = totalRows;
+ }
+
public boolean isValid() {
return valid;
}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeResource.java
index 61cf4ed204..2f26cd2f5a 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeResource.java
@@ -336,22 +336,22 @@ public class DataLakeResource extends
AbstractDataLakeResource {
}
}
- @PostMapping(path = "/measurements/{measurementID}", produces =
MediaType.APPLICATION_JSON_VALUE, consumes = MediaType.APPLICATION_JSON_VALUE)
- @PreAuthorize("this.hasWriteAuthority()")
+ @PostMapping(path = "/measurements/{measureName}", produces =
MediaType.APPLICATION_JSON_VALUE, consumes = MediaType.APPLICATION_JSON_VALUE)
+ @PreAuthorize("this.hasWriteAuthority() and
this.checkPermissionByName(#measureName, 'WRITE')")
@Operation(summary = "Store a measurement series to a data lake with the
given id", tags = {
"Data Lake" }, responses = {
@ApiResponse(responseCode = "400", description = "Can't store the
given data to this data lake"),
@ApiResponse(responseCode = "200", description = "Successfully
stored data") })
public ResponseEntity<?> storeDataToMeasurement(
- @PathVariable String measurementID,
+ @PathVariable String measureName,
@RequestBody SpQueryResult queryResult,
@Parameter(in = ParameterIn.QUERY, description = "should not identical
schemas be stored") @RequestParam(value = "ignoreSchemaMismatch", required =
false) boolean ignoreSchemaMismatch) {
var dataWriter = new DataLakeDataWriter(ignoreSchemaMismatch,
datasetStorage);
try {
- dataWriter.writeData(measurementID, queryResult);
+ dataWriter.writeData(measureName, queryResult);
} catch (SpRuntimeException e) {
LOG.warn("Could not store event", e);
- return badRequest(Notifications.error("Could not store event for
measurement " + measurementID, e.getMessage()));
+ return badRequest(Notifications.error("Could not store event for
measurement " + measureName, e.getMessage()));
}
return ok();
}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/CsvDataLakeImportService.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/CsvDataLakeImportService.java
index d68f4ec6fb..e3eb60f998 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/CsvDataLakeImportService.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/CsvDataLakeImportService.java
@@ -21,6 +21,8 @@ package org.apache.streampipes.rest.impl.datalake.importer;
import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement;
import org.apache.streampipes.model.datalake.DataLakeMeasure;
import org.apache.streampipes.model.datalake.importer.CsvImportColumn;
+import org.apache.streampipes.model.datalake.importer.CsvImportJobStartResult;
+import org.apache.streampipes.model.datalake.importer.CsvImportJobStatus;
import org.apache.streampipes.model.datalake.importer.CsvImportPreviewRequest;
import org.apache.streampipes.model.datalake.importer.CsvImportPreviewResult;
import org.apache.streampipes.model.datalake.importer.CsvImportRequest;
@@ -46,6 +48,7 @@ import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.function.Function;
+import java.util.function.IntConsumer;
import java.util.stream.Collectors;
public class CsvDataLakeImportService {
@@ -59,6 +62,7 @@ public class CsvDataLakeImportService {
private final CsvImportParser parser;
private final CsvImportValidationService validationService;
private final IDataExplorerSchemaManagement schemaManagement;
+ private final CsvImportJobManager jobManager;
public CsvDataLakeImportService(IDataExplorerSchemaManagement
schemaManagement,
IDataLakeMeasureStorage datasetStorage) {
@@ -92,13 +96,14 @@ public class CsvDataLakeImportService {
this.uploadStorage = uploadStorage;
this.parser = parser;
this.validationService = new CsvImportValidationService(schemaManagement);
+ this.jobManager = new CsvImportJobManager(this);
}
public CsvImportPreviewResult preview(CsvImportPreviewRequest request) {
var validationMessages = validationService.validatePreviewRequest(request);
var headers = parser.sanitizeHeaders(request.getHeaders());
var rows =
Optional.ofNullable(request.getRows()).orElseGet(Collections::emptyList);
- return buildPreviewResult(request, headers, rows, validationMessages,
null);
+ return buildPreviewResult(request, headers, rows, rows.size(),
validationMessages, null);
}
public CsvImportPreviewResult preview(CsvImportPreviewRequest request,
String principalSid) {
@@ -114,7 +119,14 @@ public class CsvDataLakeImportService {
try {
var upload = resolveUpload(request.getUploadId(), principalSid);
var csvSample = parser.readCsvSample(upload.path(),
request.getCsvConfig(), MAX_ANALYSIS_ROWS);
- return buildPreviewResult(request, csvSample.headers(),
csvSample.rows(), validationMessages, upload.uploadId());
+ uploadStorage.updateTotalRows(upload.uploadId(), csvSample.totalRows());
+ return buildPreviewResult(
+ request,
+ csvSample.headers(),
+ csvSample.rows(),
+ csvSample.totalRows(),
+ validationMessages,
+ upload.uploadId());
} catch (CsvImportValidationException e) {
return buildInvalidPreviewResult(e.getValidationMessages(),
request.getUploadId());
} catch (IOException | UncheckedIOException e) {
@@ -134,7 +146,14 @@ public class CsvDataLakeImportService {
var upload = uploadStorage.store(file, principalSid);
try {
var csvSample = parser.readCsvSample(upload.path(),
request.getCsvConfig(), MAX_ANALYSIS_ROWS);
- return buildPreviewResult(request, csvSample.headers(),
csvSample.rows(), validationMessages, upload.uploadId());
+ uploadStorage.updateTotalRows(upload.uploadId(), csvSample.totalRows());
+ return buildPreviewResult(
+ request,
+ csvSample.headers(),
+ csvSample.rows(),
+ csvSample.totalRows(),
+ validationMessages,
+ upload.uploadId());
} catch (IOException | UncheckedIOException e) {
uploadStorage.remove(upload.uploadId());
throw e;
@@ -162,8 +181,13 @@ public class CsvDataLakeImportService {
}
public CsvImportResult importData(CsvImportRequest request, String
principalSid) {
+ return importData(request, principalSid, importedRows -> {
+ });
+ }
+
+ CsvImportResult importData(CsvImportRequest request, String principalSid,
IntConsumer progressConsumer) {
if (hasUploadId(request)) {
- return importUploadedData(request, principalSid);
+ return importUploadedData(request, principalSid, progressConsumer);
}
var validationMessages =
validationService.validateInlineImportRequest(request);
@@ -189,11 +213,29 @@ public class CsvDataLakeImportService {
measure.measure(),
request.getColumns().stream().map(CsvImportColumn::getRuntimeName).collect(Collectors.toList()),
rows);
+ progressConsumer.accept(rows.size());
return buildImportResult(measure, rows.size());
}
- private CsvImportResult importUploadedData(CsvImportRequest request, String
principalSid) {
+ public CsvImportJobStartResult startImportJob(CsvImportRequest request,
String principalSid) {
+ return jobManager.start(request, principalSid, getTotalRows(request,
principalSid));
+ }
+
+ public Optional<CsvImportJobStatus> getImportJobStatus(String jobId, String
principalSid) {
+ return jobManager.getStatus(jobId, principalSid);
+ }
+
+ void cleanupUpload(CsvImportRequest request) {
+ if (hasUploadId(request)) {
+ uploadStorage.remove(request.getUploadId());
+ }
+ }
+
+ private CsvImportResult importUploadedData(
+ CsvImportRequest request,
+ String principalSid,
+ IntConsumer progressConsumer) {
var validationMessages =
validationService.validateStoredImportRequest(request);
if (!validationMessages.isEmpty()) {
throw new CsvImportValidationException(validationMessages);
@@ -212,12 +254,18 @@ public class CsvDataLakeImportService {
var measure = resolveTargetMeasurement(request, principalSid, eventSchema);
try {
- var importedRowCount = parser.importCsvFile(upload.path(), request,
measure.measure(), dataWriter);
- uploadStorage.remove(upload.uploadId());
+ var importedRowCount = parser.importCsvFile(
+ upload.path(),
+ request,
+ measure.measure(),
+ dataWriter,
+ progressConsumer);
return buildImportResult(measure, importedRowCount);
} catch (IOException | UncheckedIOException e) {
throw new CsvImportValidationException(List.of(
message("file", "The CSV file could not be parsed with the current
settings.")));
+ } finally {
+ uploadStorage.remove(upload.uploadId());
}
}
@@ -225,6 +273,7 @@ public class CsvDataLakeImportService {
CsvImportPreviewRequest request,
List<String> headers,
List<List<String>> rows,
+ int totalRows,
List<CsvImportValidationMessage> validationMessages,
String uploadId) {
var messages = new ArrayList<>(validationMessages);
@@ -239,6 +288,7 @@ public class CsvDataLakeImportService {
result.setPreviewRows(rows.stream().limit(MAX_PREVIEW_ROWS).collect(Collectors.toList()));
result.setColumns(columns);
result.setGuessedEventSchema(eventSchema);
+ result.setTotalRows(totalRows);
result.setTimestampCandidates(columns.stream()
.filter(CsvImportColumn::isTimestampCandidate)
.map(CsvImportColumn::getRuntimeName)
@@ -258,6 +308,13 @@ public class CsvDataLakeImportService {
return result;
}
+ private int getTotalRows(CsvImportRequest request, String principalSid) {
+ if (hasUploadId(request)) {
+ return resolveUpload(request.getUploadId(), principalSid).totalRows();
+ }
+ return
Optional.ofNullable(request.getRows()).orElseGet(Collections::emptyList).size();
+ }
+
private CsvImportUploadStorage.StoredUpload resolveUpload(String uploadId,
String principalSid) {
var upload = uploadStorage.get(uploadId).orElseThrow(() -> new
CsvImportValidationException(List.of(
message("uploadId", "The uploaded CSV file was not found. Please
upload the file again."))));
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/CsvImportJobManager.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/CsvImportJobManager.java
new file mode 100644
index 0000000000..c7f5600108
--- /dev/null
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/CsvImportJobManager.java
@@ -0,0 +1,162 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.rest.impl.datalake.importer;
+
+import org.apache.streampipes.model.datalake.importer.CsvImportJobStartResult;
+import org.apache.streampipes.model.datalake.importer.CsvImportJobState;
+import org.apache.streampipes.model.datalake.importer.CsvImportJobStatus;
+import org.apache.streampipes.model.datalake.importer.CsvImportRequest;
+import org.apache.streampipes.model.datalake.importer.CsvImportResult;
+import
org.apache.streampipes.model.datalake.importer.CsvImportValidationMessage;
+
+import java.time.Duration;
+import java.time.Instant;
+import java.util.List;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.UUID;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+
+class CsvImportJobManager {
+
+ private static final Duration DEFAULT_TTL = Duration.ofHours(12);
+
+ private final CsvDataLakeImportService importService;
+ private final ConcurrentMap<String, StoredJob> jobs;
+ private final ExecutorService executorService;
+ private final Duration ttl;
+
+ CsvImportJobManager(CsvDataLakeImportService importService) {
+ this(importService, Executors.newCachedThreadPool(), DEFAULT_TTL);
+ }
+
+ CsvImportJobManager(
+ CsvDataLakeImportService importService,
+ ExecutorService executorService,
+ Duration ttl) {
+ this.importService = importService;
+ this.executorService = executorService;
+ this.ttl = ttl;
+ this.jobs = new ConcurrentHashMap<>();
+ }
+
+ CsvImportJobStartResult start(CsvImportRequest request, String ownerSid, int
totalRows) {
+ cleanupExpired();
+
+ var jobId = UUID.randomUUID().toString();
+ var job = new StoredJob(jobId, ownerSid, totalRows);
+ jobs.put(jobId, job);
+
+ executorService.submit(() -> runImport(job, request, ownerSid));
+ return new CsvImportJobStartResult(jobId);
+ }
+
+ Optional<CsvImportJobStatus> getStatus(String jobId, String ownerSid) {
+ cleanupExpired();
+
+ return Optional.ofNullable(jobs.get(jobId))
+ .filter(job -> Objects.equals(job.ownerSid(), ownerSid))
+ .map(StoredJob::status);
+ }
+
+ private void runImport(StoredJob job, CsvImportRequest request, String
ownerSid) {
+ try {
+ var result = importService.importData(request, ownerSid,
job::setProcessedRows);
+ job.succeeded(result);
+ } catch (CsvImportValidationException e) {
+ job.failed(e.getValidationMessages());
+ } catch (RuntimeException e) {
+ job.failed(List.of(new CsvImportValidationMessage("import", "CSV import
failed.")));
+ } finally {
+ importService.cleanupUpload(request);
+ }
+ }
+
+ private void cleanupExpired() {
+ var expiresBefore = Instant.now().minus(ttl);
+ jobs.entrySet().removeIf(entry ->
entry.getValue().createdAt().isBefore(expiresBefore));
+ }
+
+ private static final class StoredJob {
+
+ private final String jobId;
+ private final String ownerSid;
+ private final int totalRows;
+ private final Instant createdAt;
+ private CsvImportJobState state;
+ private int processedRows;
+ private CsvImportResult result;
+ private List<CsvImportValidationMessage> validationMessages;
+
+ private StoredJob(String jobId, String ownerSid, int totalRows) {
+ this.jobId = jobId;
+ this.ownerSid = ownerSid;
+ this.totalRows = totalRows;
+ this.createdAt = Instant.now();
+ this.state = CsvImportJobState.RUNNING;
+ this.validationMessages = List.of();
+ }
+
+ synchronized CsvImportJobStatus status() {
+ var status = new CsvImportJobStatus();
+ status.setJobId(jobId);
+ status.setState(state);
+ status.setProcessedRows(processedRows);
+ status.setTotalRows(totalRows);
+ status.setProgress(calculateProgress());
+ status.setResult(result);
+ status.setValidationMessages(validationMessages);
+ return status;
+ }
+
+ synchronized void setProcessedRows(int processedRows) {
+ this.processedRows = Math.min(processedRows, totalRows);
+ }
+
+ synchronized void succeeded(CsvImportResult result) {
+ this.result = result;
+ this.processedRows = totalRows > 0 ? totalRows :
result.getImportedRowCount();
+ this.state = CsvImportJobState.SUCCEEDED;
+ this.validationMessages = List.of();
+ }
+
+ synchronized void failed(List<CsvImportValidationMessage>
validationMessages) {
+ this.state = CsvImportJobState.FAILED;
+ this.validationMessages = validationMessages;
+ }
+
+ String ownerSid() {
+ return ownerSid;
+ }
+
+ Instant createdAt() {
+ return createdAt;
+ }
+
+ private int calculateProgress() {
+ if (totalRows <= 0) {
+ return state == CsvImportJobState.SUCCEEDED ? 100 : 0;
+ }
+ return Math.min(100, Math.round((processedRows * 100.0f) / totalRows));
+ }
+ }
+}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/CsvImportParser.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/CsvImportParser.java
index ee43045a4a..772ff87691 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/CsvImportParser.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/CsvImportParser.java
@@ -48,6 +48,7 @@ import java.util.List;
import java.util.Locale;
import java.util.Objects;
import java.util.Optional;
+import java.util.function.IntConsumer;
import java.util.stream.Collectors;
class CsvImportParser {
@@ -153,6 +154,17 @@ class CsvImportParser {
CsvImportRequest request,
DataLakeMeasure measure,
DataLakeDataWriter dataWriter
+ ) throws IOException {
+ return importCsvFile(path, request, measure, dataWriter, importedRows -> {
+ });
+ }
+
+ int importCsvFile(
+ Path path,
+ CsvImportRequest request,
+ DataLakeMeasure measure,
+ DataLakeDataWriter dataWriter,
+ IntConsumer progressConsumer
) throws IOException {
var runtimeHeaders = request.getColumns().stream()
.map(CsvImportColumn::getRuntimeName)
@@ -178,8 +190,10 @@ class CsvImportParser {
}
batch.add(convertRow(row, request, rowNumber));
if (batch.size() >= IMPORT_BATCH_SIZE) {
+ var batchSize = batch.size();
dataWriter.writeData(measure, runtimeHeaders, new
ArrayList<>(batch));
- importedRows[0] += IMPORT_BATCH_SIZE;
+ importedRows[0] += batchSize;
+ progressConsumer.accept(importedRows[0]);
batch.clear();
}
}
@@ -189,6 +203,7 @@ class CsvImportParser {
var batchSize = batch.size();
dataWriter.writeData(measure, runtimeHeaders, new ArrayList<>(batch));
importedRows[0] += batchSize;
+ progressConsumer.accept(importedRows[0]);
}
return importedRows[0];
@@ -197,6 +212,7 @@ class CsvImportParser {
CsvFileSample readCsvSample(Path path, CsvImportConfiguration config, int
maxRows) throws IOException {
var headers = new ArrayList<String>();
var rows = new ArrayList<List<String>>();
+ var totalRows = new int[]{0};
parseCsvFile(path, config, new CsvRowConsumer() {
@Override
@@ -211,6 +227,7 @@ class CsvImportParser {
message("rows", "Row " + rowNumber + " does not match the header
size.")
));
}
+ totalRows[0] += 1;
if (rows.size() < maxRows) {
rows.add(row);
}
@@ -224,7 +241,7 @@ class CsvImportParser {
throw new CsvImportValidationException(List.of(message("rows", "At least
one row must be provided.")));
}
- return new CsvFileSample(headers, rows);
+ return new CsvFileSample(headers, rows, totalRows[0]);
}
private List<Object> convertRow(List<String> row, CsvImportRequest request,
int rowNumber) {
@@ -602,10 +619,12 @@ class CsvImportParser {
private final List<String> headers;
private final List<List<String>> rows;
+ private final int totalRows;
- CsvFileSample(List<String> headers, List<List<String>> rows) {
+ CsvFileSample(List<String> headers, List<List<String>> rows, int
totalRows) {
this.headers = headers;
this.rows = rows;
+ this.totalRows = totalRows;
}
List<String> headers() {
@@ -615,5 +634,9 @@ class CsvImportParser {
List<List<String>> rows() {
return rows;
}
+
+ int totalRows() {
+ return totalRows;
+ }
}
}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/CsvImportUploadStorage.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/CsvImportUploadStorage.java
index a62dc1495c..7c78fb2f23 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/CsvImportUploadStorage.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/CsvImportUploadStorage.java
@@ -60,6 +60,7 @@ class CsvImportUploadStorage {
uploadId,
tempFile,
ownerSid,
+ 0,
Instant.now()
);
uploads.put(uploadId, upload);
@@ -71,6 +72,13 @@ class CsvImportUploadStorage {
return Optional.ofNullable(uploads.get(uploadId));
}
+ void updateTotalRows(String uploadId, int totalRows) {
+ var upload = uploads.get(uploadId);
+ if (upload != null) {
+ upload.setTotalRows(totalRows);
+ }
+ }
+
void remove(String uploadId) {
var removed = uploads.remove(uploadId);
if (removed != null) {
@@ -102,12 +110,14 @@ class CsvImportUploadStorage {
private final String uploadId;
private final Path path;
private final String ownerSid;
+ private int totalRows;
private final Instant createdAt;
- StoredUpload(String uploadId, Path path, String ownerSid, Instant
createdAt) {
+ StoredUpload(String uploadId, Path path, String ownerSid, int totalRows,
Instant createdAt) {
this.uploadId = uploadId;
this.path = path;
this.ownerSid = ownerSid;
+ this.totalRows = totalRows;
this.createdAt = createdAt;
}
@@ -123,6 +133,14 @@ class CsvImportUploadStorage {
return ownerSid;
}
+ int totalRows() {
+ return totalRows;
+ }
+
+ void setTotalRows(int totalRows) {
+ this.totalRows = totalRows;
+ }
+
Instant createdAt() {
return createdAt;
}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/DataLakeImportResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/DataLakeImportResource.java
index 2ad1866bf3..1538aaae1b 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/DataLakeImportResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/DataLakeImportResource.java
@@ -19,6 +19,7 @@
package org.apache.streampipes.rest.impl.datalake.importer;
import
org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator;
+import org.apache.streampipes.model.datalake.importer.CsvImportJobStatus;
import org.apache.streampipes.model.datalake.importer.CsvImportPreviewRequest;
import org.apache.streampipes.model.datalake.importer.CsvImportPreviewResult;
import org.apache.streampipes.model.datalake.importer.CsvImportRequest;
@@ -34,6 +35,8 @@ import org.springframework.http.HttpStatus;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
import org.springframework.security.access.prepost.PreAuthorize;
+import org.springframework.web.bind.annotation.GetMapping;
+import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
@@ -129,13 +132,13 @@ public class DataLakeImportResource extends
AbstractDataLakeResource {
produces = MediaType.APPLICATION_JSON_VALUE
)
@PreAuthorize("this.hasWriteAuthority()")
- public ResponseEntity<CsvImportResult> importData(@RequestBody
CsvImportRequest request) {
+ public ResponseEntity<?> importData(@RequestBody CsvImportRequest request) {
if (!hasWritePermission(request.getTarget())) {
return ResponseEntity.status(HttpStatus.FORBIDDEN).build();
}
try {
- return ok(importService.importData(request, getAuthenticatedUserSid()));
+ return ok(importService.startImportJob(request,
getAuthenticatedUserSid()));
} catch (CsvImportValidationException e) {
var result = new CsvImportResult();
result.setValidationMessages(e.getValidationMessages());
@@ -143,6 +146,20 @@ public class DataLakeImportResource extends
AbstractDataLakeResource {
}
}
+ /**
+ * Returns the current state of a CSV import job owned by the authenticated
user.
+ */
+ @GetMapping(
+ path = "/{jobId}",
+ produces = MediaType.APPLICATION_JSON_VALUE
+ )
+ @PreAuthorize("this.hasWriteAuthority()")
+ public ResponseEntity<CsvImportJobStatus> getImportStatus(@PathVariable
String jobId) {
+ return importService.getImportJobStatus(jobId, getAuthenticatedUserSid())
+ .map(ResponseEntity::ok)
+ .orElseGet(() -> ResponseEntity.notFound().build());
+ }
+
private boolean
hasWritePermission(org.apache.streampipes.model.datalake.importer.CsvImportTarget
target) {
return target == null
|| target.getMode() != CsvImportTargetMode.EXISTING
diff --git
a/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/DataLakeResourceTest.java
b/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/DataLakeResourceTest.java
index aea4c1e753..c06c2e6337 100644
---
a/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/DataLakeResourceTest.java
+++
b/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/DataLakeResourceTest.java
@@ -19,9 +19,11 @@
package org.apache.streampipes.rest.impl.datalake;
import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement;
+import org.apache.streampipes.model.datalake.SpQueryResult;
import org.junit.jupiter.api.Test;
import org.springframework.http.HttpStatus;
+import org.springframework.security.access.prepost.PreAuthorize;
import java.lang.reflect.Field;
import java.util.List;
@@ -39,6 +41,21 @@ import static org.mockito.Mockito.when;
class DataLakeResourceTest {
+ @Test
+ void storeDataToMeasurementRequiresWritePermissionForMeasurement() throws
Exception {
+ var method = DataLakeResource.class.getMethod(
+ "storeDataToMeasurement",
+ String.class,
+ SpQueryResult.class,
+ boolean.class);
+
+ var annotation = method.getAnnotation(PreAuthorize.class);
+
+ assertEquals(
+ "this.hasWriteAuthority() and this.checkPermissionByName(#measureName,
'WRITE')",
+ annotation.value());
+ }
+
@Test
void getLatestEventsRequestsLatestTimestampsForDistinctMeasurements() throws
Exception {
var queryManagement = mock(IDataExplorerQueryManagement.class);
diff --git
a/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/importer/CsvDataLakeImportServiceTest.java
b/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/importer/CsvDataLakeImportServiceTest.java
index 5cd851be5a..508f542dfd 100644
---
a/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/importer/CsvDataLakeImportServiceTest.java
+++
b/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/importer/CsvDataLakeImportServiceTest.java
@@ -22,6 +22,8 @@ import
org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement;
import org.apache.streampipes.model.datalake.DataLakeMeasure;
import org.apache.streampipes.model.datalake.importer.CsvImportColumn;
import org.apache.streampipes.model.datalake.importer.CsvImportConfiguration;
+import org.apache.streampipes.model.datalake.importer.CsvImportJobState;
+import org.apache.streampipes.model.datalake.importer.CsvImportJobStatus;
import org.apache.streampipes.model.datalake.importer.CsvImportPreviewRequest;
import org.apache.streampipes.model.datalake.importer.CsvImportRequest;
import org.apache.streampipes.model.datalake.importer.CsvImportSchemaIssueType;
@@ -40,6 +42,7 @@ import org.springframework.web.multipart.MultipartFile;
import java.io.ByteArrayInputStream;
import java.util.List;
import java.util.Optional;
+import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
@@ -68,6 +71,7 @@ class CsvDataLakeImportServiceTest {
var result = service.preview(makePreviewRequest("existing-measure"));
assertFalse(result.isValid());
+ assertEquals(2, result.getTotalRows());
assertEquals("LONG", result.getColumns().get(0).getInferredType());
assertEquals("FLOAT", result.getColumns().get(1).getInferredType());
assertTrue(result.getValidationMessages()
@@ -181,6 +185,7 @@ class CsvDataLakeImportServiceTest {
assertTrue(previewResult.isValid());
assertEquals(2, previewResult.getPreviewRows().size());
+ assertEquals(2, previewResult.getTotalRows());
assertTrue(previewResult.getUploadId() != null &&
!previewResult.getUploadId().isBlank());
var importRequest = new CsvImportRequest();
@@ -197,6 +202,50 @@ class CsvDataLakeImportServiceTest {
verify(dataWriter).writeData(any(DataLakeMeasure.class), anyList(),
anyList());
}
+ @Test
+ void shouldStartImportJobAndExposeSucceededStatus() throws Exception {
+ var schemaManagement = mock(IDataExplorerSchemaManagement.class);
+ var dataWriter = mock(DataLakeDataWriter.class);
+ var service = new CsvDataLakeImportService(schemaManagement, dataWriter);
+
+ when(schemaManagement.getExistingMeasureByName("new-measure"))
+ .thenReturn(Optional.empty());
+
when(schemaManagement.createOrUpdateMeasurement(any(DataLakeMeasure.class),
any()))
+ .thenAnswer(invocation -> {
+ var measure = invocation.getArgument(0, DataLakeMeasure.class);
+ measure.setElementId("measure-id");
+ return measure;
+ });
+
+ var startResult =
service.startImportJob(makeImportRequest(CsvImportTargetMode.NEW,
"new-measure"), "sid");
+ var status = awaitTerminalStatus(service, startResult.getJobId(), "sid");
+
+ assertEquals(CsvImportJobState.SUCCEEDED, status.getState());
+ assertEquals(2, status.getProcessedRows());
+ assertEquals(2, status.getTotalRows());
+ assertEquals(100, status.getProgress());
+ assertEquals(2, status.getResult().getImportedRowCount());
+ assertTrue(service.getImportJobStatus(startResult.getJobId(),
"other-sid").isEmpty());
+ }
+
+ @Test
+ void shouldExposeFailedImportJobValidationMessages() throws Exception {
+ var schemaManagement = mock(IDataExplorerSchemaManagement.class);
+ var dataWriter = mock(DataLakeDataWriter.class);
+ var service = new CsvDataLakeImportService(schemaManagement, dataWriter);
+
+ var request = makeImportRequest(CsvImportTargetMode.NEW, "new-measure");
+ request.setTimestampColumn(null);
+
+ var startResult = service.startImportJob(request, "sid");
+ var status = awaitTerminalStatus(service, startResult.getJobId(), "sid");
+
+ assertEquals(CsvImportJobState.FAILED, status.getState());
+ assertTrue(status.getValidationMessages()
+ .stream()
+ .anyMatch(message -> message.getMessage().contains("timestamp")));
+ }
+
@Test
void shouldRejectMissingTimestampValuesInUploadedCsv() throws Exception {
var schemaManagement = mock(IDataExplorerSchemaManagement.class);
@@ -367,6 +416,22 @@ class CsvDataLakeImportServiceTest {
return new EventSchema(List.of(timestamp, temperature));
}
+ private CsvImportJobStatus awaitTerminalStatus(
+ CsvDataLakeImportService service,
+ String jobId,
+ String sid
+ ) throws Exception {
+ CsvImportJobStatus status = null;
+ for (int i = 0; i < 50; i++) {
+ status = service.getImportJobStatus(jobId, sid).orElseThrow();
+ if (status.getState() != CsvImportJobState.RUNNING) {
+ return status;
+ }
+ TimeUnit.MILLISECONDS.sleep(20);
+ }
+ return status;
+ }
+
private EventSchema makeStoredExistingSchema() {
var temperature = new EventPropertyPrimitive();
temperature.setRuntimeName("temperature");
diff --git
a/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/importer/CsvImportParserTest.java
b/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/importer/CsvImportParserTest.java
index 5857f4fe5c..709f50b6f8 100644
---
a/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/importer/CsvImportParserTest.java
+++
b/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/importer/CsvImportParserTest.java
@@ -69,6 +69,7 @@ class CsvImportParserTest {
assertEquals(List.of("timestamp", "text"), sample.headers());
assertEquals("a, b", sample.rows().get(0).get(1));
assertEquals("escaped \"quote\"", sample.rows().get(1).get(1));
+ assertEquals(2, sample.totalRows());
}
@Test
diff --git a/ui/STYLEGUIDE.md b/ui/STYLEGUIDE.md
index 752a40a7df..59574173fb 100644
--- a/ui/STYLEGUIDE.md
+++ b/ui/STYLEGUIDE.md
@@ -126,6 +126,20 @@ Use it as follows:
Allowed types are `info`, `warning`, `error` and `success`.
You can also add additional content to the banner.
+#### Progress bar
+
+Use `sp-progress-bar` when real progress is available. Provide the current
`value`, the `max` value, and optional title or item label text.
+
+```html
+<sp-progress-bar
+ [title]="'Uploading CSV data' | translate"
+ [ariaLabel]="'CSV import progress' | translate"
+ [value]="processedRows"
+ [max]="totalRows"
+ [itemLabel]="'rows imported' | translate"
+></sp-progress-bar>
+```
+
#### Tables
For rendering tables, always use the `sp-table` component which comes with
pre-defined features for paging, sorting and layout.
diff --git a/ui/cypress/support/utils/dataset/DataLakeSeedUtils.ts
b/ui/cypress/support/utils/dataset/DataLakeSeedUtils.ts
index a1806d6c68..4d3153ee9c 100644
--- a/ui/cypress/support/utils/dataset/DataLakeSeedUtils.ts
+++ b/ui/cypress/support/utils/dataset/DataLakeSeedUtils.ts
@@ -20,6 +20,7 @@ import * as CSV from 'csv-string';
type CsvRuntimeType = 'STRING' | 'BOOLEAN' | 'LONG' | 'FLOAT';
type CsvImportTargetMode = 'NEW' | 'EXISTING';
+type CsvImportJobState = 'RUNNING' | 'SUCCEEDED' | 'FAILED';
interface CsvImportConfiguration {
delimiter: string;
@@ -53,6 +54,19 @@ interface CsvImportResult {
validationMessages: Array<{ field: string; message: string }>;
}
+interface CsvImportJobStartResult {
+ jobId: string;
+}
+
+interface CsvImportJobStatus {
+ state: CsvImportJobState;
+ processedRows: number;
+ totalRows: number;
+ progress: number;
+ result?: CsvImportResult;
+ validationMessages: Array<{ field: string; message: string }>;
+}
+
interface ColumnOverride {
runtimeName?: string;
runtimeType?: CsvRuntimeType;
@@ -240,7 +254,7 @@ export class DataLakeSeedUtils {
};
return cy
- .request<CsvImportResult>({
+ .request<CsvImportJobStartResult>({
method: 'POST',
url: '/streampipes-backend/api/v4/datalake/import',
body: request,
@@ -248,21 +262,66 @@ export class DataLakeSeedUtils {
Authorization: `Bearer ${token}`,
},
})
- .then(importResponse => {
+ .then(importResponse =>
+ this.pollImportStatus(
+ importResponse.body.jobId,
+ token,
+ ),
+ )
+ .then(importResult => {
expect(
- importResponse.body.validationMessages,
+ importResult.validationMessages,
'import validation messages',
).to.have.length(0);
expect(
- importResponse.body.importedRowCount,
+ importResult.importedRowCount,
'imported row count',
).to.equal(options.rows.length);
- return importResponse.body;
+ return importResult;
});
});
});
}
+ private static pollImportStatus(
+ jobId: string,
+ token: string | null,
+ retries = 60,
+ ): Cypress.Chainable<CsvImportResult> {
+ return cy
+ .request<CsvImportJobStatus>({
+ method: 'GET',
+ url:
`/streampipes-backend/api/v4/datalake/import/${encodeURIComponent(jobId)}`,
+ headers: {
+ Authorization: `Bearer ${token}`,
+ },
+ })
+ .then(statusResponse => {
+ const status = statusResponse.body;
+ if (status.state === 'SUCCEEDED' && status.result) {
+ return status.result;
+ }
+
+ if (status.state === 'FAILED') {
+ throw new Error(
+ status.validationMessages
+ .map(message => message.message)
+ .join('\n') || 'CSV import failed.',
+ );
+ }
+
+ if (retries <= 0) {
+ throw new Error('CSV import did not finish in time.');
+ }
+
+ return cy
+ .wait(1000)
+ .then(() =>
+ this.pollImportStatus(jobId, token, retries - 1),
+ );
+ });
+ }
+
private static buildColumns(
previewColumns: CsvImportColumn[],
timestampColumn: string,
diff --git a/ui/deployment/i18n/de.json b/ui/deployment/i18n/de.json
index 9a5288c723..f132229e81 100644
--- a/ui/deployment/i18n/de.json
+++ b/ui/deployment/i18n/de.json
@@ -257,6 +257,8 @@
"Creating adapter {{adapterName}}": "Adapter erstellen {{adapterName}}",
"Creating pipeline to persist data stream": "Erstellen einer Pipeline zum
Speichern der Daten",
"Critical datatype changes": "Kritische Datentypänderungen",
+ "CSV import progress": "Fortschritt des CSV-Imports",
+ "CSV import status could not be loaded.": "Der CSV-Importstatus konnte nicht
geladen werden.",
"Current / Target": "Aktuell / Ziel",
"Current Warning Range: ": "Aktueller Warnbereich: ",
"Current day": "Aktueller Tag",
@@ -1114,7 +1116,6 @@
"The current data selection can't be displayed by this chart.": "Die
aktuelle Auswahl kann in diesem Diagramm nicht angezeigt werden.",
"The current value displayed as a number": "Der aktuelle Wert wird als Zahl
angezeigt",
"The current value displayed in a gauge": "Der aktuelle Wert, der im
Gauge-Chart angezeigt wird",
- "The data is currently being written to the data lake.": "Die Daten werden
derzeit in den Data Lake geschrieben.",
"The data type of the field values": "Der Datentyp der Feldwerte",
"The default logo": "Das Standardlogo",
"The desired adapter was not found!": "Der gewünschte Adapter wurde nicht
gefunden!",
diff --git a/ui/deployment/i18n/en.json b/ui/deployment/i18n/en.json
index a2a0063823..675ee0f3a1 100644
--- a/ui/deployment/i18n/en.json
+++ b/ui/deployment/i18n/en.json
@@ -257,6 +257,8 @@
"Creating adapter {{adapterName}}": "Creating adapter {{adapterName}}",
"Creating pipeline to persist data stream": null,
"Critical datatype changes": null,
+ "CSV import progress": null,
+ "CSV import status could not be loaded.": null,
"Current / Target": null,
"Current Warning Range: ": null,
"Current day": null,
@@ -1114,7 +1116,6 @@
"The current data selection can't be displayed by this chart.": null,
"The current value displayed as a number": null,
"The current value displayed in a gauge": null,
- "The data is currently being written to the data lake.": null,
"The data type of the field values": null,
"The default logo": null,
"The desired adapter was not found!": null,
diff --git a/ui/deployment/i18n/pl.json b/ui/deployment/i18n/pl.json
index a54c8b4a4a..cc153b2023 100644
--- a/ui/deployment/i18n/pl.json
+++ b/ui/deployment/i18n/pl.json
@@ -257,6 +257,8 @@
"Creating adapter {{adapterName}}": "Tworzenie adaptera {{adapterName}}",
"Creating pipeline to persist data stream": "Tworzenie strumienia do
utrwalenia strumienia danych",
"Critical datatype changes": "Krytyczne zmiany typów danych",
+ "CSV import progress": "Postęp importu CSV",
+ "CSV import status could not be loaded.": "Nie można wczytać statusu importu
CSV.",
"Current / Target": "Bieżąca / Docelowa",
"Current Warning Range: ": "Bieżący zakres ostrzegawczy: ",
"Current day": "Bieżący dzień",
@@ -1114,7 +1116,6 @@
"The current data selection can't be displayed by this chart.": "Bieżącego
wyboru danych nie można wyświetlić na tym wykresie.",
"The current value displayed as a number": "Bieżąca wartość wyświetlana jako
liczba",
"The current value displayed in a gauge": "Bieżąca wartość wyświetlana na
wskaźniku",
- "The data is currently being written to the data lake.": "Dane są obecnie
zapisywane do jeziora danych.",
"The data type of the field values": "Typ danych wartości pola",
"The default logo": "Domyślne logo",
"The desired adapter was not found!": "Nie znaleziono żądanego adaptera!",
diff --git
a/ui/projects/streampipes/platform-services/src/lib/apis/datalake-rest.service.ts
b/ui/projects/streampipes/platform-services/src/lib/apis/datalake-rest.service.ts
index ff7769a800..aa4c4d7680 100644
---
a/ui/projects/streampipes/platform-services/src/lib/apis/datalake-rest.service.ts
+++
b/ui/projects/streampipes/platform-services/src/lib/apis/datalake-rest.service.ts
@@ -34,10 +34,11 @@ import {
ResourceSummaryDto,
} from '../model/resource/resource-summary.model';
import {
+ CsvImportJobStartResult,
+ CsvImportJobStatus,
CsvImportPreviewRequest,
CsvImportPreviewResult,
CsvImportRequest,
- CsvImportResult,
CsvImportSchemaValidationRequest,
CsvImportSchemaValidationResult,
} from '../model/datalake/csv-import.model';
@@ -310,8 +311,19 @@ export class DatalakeRestService {
);
}
- importCsvData(request: CsvImportRequest): Observable<CsvImportResult> {
- return this.http.post<CsvImportResult>(this.dataLakeImportUrl,
request);
+ importCsvData(
+ request: CsvImportRequest,
+ ): Observable<CsvImportJobStartResult> {
+ return this.http.post<CsvImportJobStartResult>(
+ this.dataLakeImportUrl,
+ request,
+ );
+ }
+
+ getCsvImportJobStatus(jobId: string): Observable<CsvImportJobStatus> {
+ return this.http.get<CsvImportJobStatus>(
+ `${this.dataLakeImportUrl}/${encodeURIComponent(jobId)}`,
+ );
}
dropSingleMeasurementSeries(index: string) {
diff --git
a/ui/projects/streampipes/platform-services/src/lib/model/datalake/csv-import.model.ts
b/ui/projects/streampipes/platform-services/src/lib/model/datalake/csv-import.model.ts
index 98edf75c48..2bcc5bc441 100644
---
a/ui/projects/streampipes/platform-services/src/lib/model/datalake/csv-import.model.ts
+++
b/ui/projects/streampipes/platform-services/src/lib/model/datalake/csv-import.model.ts
@@ -19,6 +19,7 @@
import { EventSchema } from '../gen/streampipes-model';
export type CsvImportTargetMode = 'NEW' | 'EXISTING';
+export type CsvImportJobState = 'RUNNING' | 'SUCCEEDED' | 'FAILED';
export type CsvRuntimeType = 'STRING' | 'BOOLEAN' | 'LONG' | 'FLOAT';
export interface CsvImportConfiguration {
@@ -66,6 +67,7 @@ export interface CsvImportPreviewResult {
columns: CsvImportColumn[];
guessedEventSchema: EventSchema;
timestampCandidates: string[];
+ totalRows: number;
valid: boolean;
validationMessages: CsvImportValidationMessage[];
}
@@ -112,3 +114,17 @@ export interface CsvImportResult {
importedRowCount: number;
validationMessages: CsvImportValidationMessage[];
}
+
+export interface CsvImportJobStartResult {
+ jobId: string;
+}
+
+export interface CsvImportJobStatus {
+ jobId: string;
+ state: CsvImportJobState;
+ processedRows: number;
+ totalRows: number;
+ progress: number;
+ result?: CsvImportResult;
+ validationMessages: CsvImportValidationMessage[];
+}
diff --git
a/ui/projects/streampipes/shared-ui/src/lib/components/progress-bar/progress-bar.component.html
b/ui/projects/streampipes/shared-ui/src/lib/components/progress-bar/progress-bar.component.html
new file mode 100644
index 0000000000..08af6e48c4
--- /dev/null
+++
b/ui/projects/streampipes/shared-ui/src/lib/components/progress-bar/progress-bar.component.html
@@ -0,0 +1,37 @@
+<!--
+ ~ Licensed to the Apache Software Foundation (ASF) under one or more
+ ~ contributor license agreements. See the NOTICE file distributed with
+ ~ this work for additional information regarding copyright ownership.
+ ~ The ASF licenses this file to You under the Apache License, Version 2.0
+ ~ (the "License"); you may not use this file except in compliance with
+ ~ the License. You may obtain a copy of the License at
+ ~
+ ~ http://www.apache.org/licenses/LICENSE-2.0
+ ~
+ ~ Unless required by applicable law or agreed to in writing, software
+ ~ distributed under the License is distributed on an "AS IS" BASIS,
+ ~ WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ ~ See the License for the specific language governing permissions and
+ ~ limitations under the License.
+ ~
+ -->
+
+<div class="progress-bar">
+ <mat-progress-bar
+ mode="determinate"
+ [value]="progressValue()"
+ [attr.aria-label]="resolvedAriaLabel()"
+ [attr.data-cy]="progressBarDataCy()"
+ class="progress-bar-track"
+ ></mat-progress-bar>
+ @if (title()) {
+ <div class="progress-bar-title">{{ title() }}</div>
+ }
+ <div class="progress-bar-copy" [attr.data-cy]="progressLabelDataCy()">
+ {{ boundedValue() | number }}
+ /
+ {{ maxValue() | number }}
+ {{ itemLabel() }}
+ </div>
+ <div class="progress-bar-meta">{{ progressValue() }}%</div>
+</div>
diff --git
a/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-upload-state/csv-import-upload-state.component.scss
b/ui/projects/streampipes/shared-ui/src/lib/components/progress-bar/progress-bar.component.scss
similarity index 67%
copy from
ui/src/app/dataset/dialog/csv-import-dialog/csv-import-upload-state/csv-import-upload-state.component.scss
copy to
ui/projects/streampipes/shared-ui/src/lib/components/progress-bar/progress-bar.component.scss
index 3b740c3d52..ebfb27990d 100644
---
a/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-upload-state/csv-import-upload-state.component.scss
+++
b/ui/projects/streampipes/shared-ui/src/lib/components/progress-bar/progress-bar.component.scss
@@ -16,43 +16,26 @@
*
*/
-.upload-state {
- min-height: 280px;
+.progress-bar {
display: flex;
flex-direction: column;
align-items: center;
- justify-content: center;
- gap: 14px;
+ gap: var(--space-md);
text-align: center;
+ width: 100%;
}
-.upload-title {
- font-size: 20px;
- font-weight: 600;
-}
-
-.upload-copy,
-.result-meta {
- color: var(--color-text-200, #5a6673);
- max-width: 420px;
+.progress-bar-track {
+ width: min(100%, 26.25rem);
}
-.success-icon {
- color: var(--color-success, #2e7d32);
- font-size: 48px;
- width: 48px;
- height: 48px;
-}
-
-.result-name {
+.progress-bar-title {
+ font-size: var(--font-size-xl);
font-weight: 600;
}
-.upload-errors {
- max-width: 520px;
+.progress-bar-copy,
+.progress-bar-meta {
color: var(--color-text-200, #5a6673);
-}
-
-.upload-error-line + .upload-error-line {
- margin-top: 6px;
+ max-width: 26.25rem;
}
diff --git
a/ui/projects/streampipes/shared-ui/src/lib/components/progress-bar/progress-bar.component.ts
b/ui/projects/streampipes/shared-ui/src/lib/components/progress-bar/progress-bar.component.ts
new file mode 100644
index 0000000000..9d00243f33
--- /dev/null
+++
b/ui/projects/streampipes/shared-ui/src/lib/components/progress-bar/progress-bar.component.ts
@@ -0,0 +1,57 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+import { DecimalPipe } from '@angular/common';
+import {
+ ChangeDetectionStrategy,
+ Component,
+ computed,
+ input,
+} from '@angular/core';
+import { MatProgressBar } from '@angular/material/progress-bar';
+
+@Component({
+ selector: 'sp-progress-bar',
+ templateUrl: './progress-bar.component.html',
+ styleUrls: ['./progress-bar.component.scss'],
+ changeDetection: ChangeDetectionStrategy.OnPush,
+ imports: [DecimalPipe, MatProgressBar],
+})
+export class ProgressBarComponent {
+ readonly value = input(0);
+ readonly max = input(100);
+ readonly title = input('');
+ readonly ariaLabel = input('');
+ readonly itemLabel = input('');
+ readonly progressBarDataCy = input<string | undefined>(undefined);
+ readonly progressLabelDataCy = input<string | undefined>(undefined);
+
+ readonly maxValue = computed(() => Math.max(0, this.max()));
+ readonly boundedValue = computed(() =>
+ Math.min(this.maxValue(), Math.max(0, this.value())),
+ );
+ readonly progressValue = computed(() => {
+ if (this.maxValue() === 0) {
+ return 0;
+ }
+ return Math.round((this.boundedValue() / this.maxValue()) * 100);
+ });
+ readonly resolvedAriaLabel = computed(
+ () => this.ariaLabel() || this.title(),
+ );
+}
diff --git a/ui/projects/streampipes/shared-ui/src/public-api.ts
b/ui/projects/streampipes/shared-ui/src/public-api.ts
index 2df6df6169..e586495e2c 100644
--- a/ui/projects/streampipes/shared-ui/src/public-api.ts
+++ b/ui/projects/streampipes/shared-ui/src/public-api.ts
@@ -46,6 +46,7 @@ export * from
'./lib/components/element-id/element-id.component';
export * from './lib/components/form-field/form-field.component';
export * from './lib/components/form-label/form-label.component';
export * from
'./lib/components/property-scope-badge/property-scope-badge.component';
+export * from './lib/components/progress-bar/progress-bar.component';
export * from './lib/components/split-section/split-section.component';
export * from './lib/components/split-button/split-button.component';
export * from
'./lib/components/sp-exception-message/sp-exception-message.component';
diff --git
a/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-dialog.component.html
b/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-dialog.component.html
index 1aef764adb..232c1c0637 100644
---
a/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-dialog.component.html
+++
b/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-dialog.component.html
@@ -259,6 +259,8 @@
[hasImportResult]="hasImportResult()"
[importResult]="importResult()"
[uploadErrors]="uploadMessages()"
+ [processedRows]="importProcessedRows()"
+ [totalRows]="importTotalRows()"
></sp-csv-import-upload-state>
</sp-split-section>
diff --git
a/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-dialog.component.ts
b/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-dialog.component.ts
index ae5773de7f..02644f9ee8 100644
--- a/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-dialog.component.ts
+++ b/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-dialog.component.ts
@@ -19,6 +19,7 @@
import {
Component,
computed,
+ DestroyRef,
inject,
Input,
signal,
@@ -39,10 +40,11 @@ import { MatInput } from '@angular/material/input';
import { MatOption, MatSelect } from '@angular/material/select';
import { MatProgressSpinner } from '@angular/material/progress-spinner';
import { MatStep, MatStepLabel, MatStepper } from '@angular/material/stepper';
-import { TranslatePipe } from '@ngx-translate/core';
+import { TranslatePipe, TranslateService } from '@ngx-translate/core';
import {
CsvImportColumn,
CsvImportConfiguration,
+ CsvImportJobStatus,
CsvImportPreviewRequest,
CsvImportPreviewResult,
CsvImportRequest,
@@ -63,7 +65,7 @@ import {
FormFieldComponent,
SplitSectionComponent,
} from '@streampipes/shared-ui';
-import { startWith } from 'rxjs';
+import { startWith, Subscription, switchMap, timer } from 'rxjs';
import { CsvImportColumnModel, CsvImportColumnRole } from './csv-import.model';
import { CsvImportPreviewTableComponent } from
'./csv-import-preview-table/csv-import-preview-table.component';
import { CsvImportUploadStateComponent } from
'./csv-import-upload-state/csv-import-upload-state.component';
@@ -102,9 +104,12 @@ export class CsvImportDialogComponent {
private readonly fb = inject(FormBuilder);
private readonly dialogRef = inject(DialogRef<CsvImportDialogComponent>);
private readonly datalakeRestService = inject(DatalakeRestService);
+ private readonly destroyRef = inject(DestroyRef);
+ private readonly translateService = inject(TranslateService);
private previewReloadTimeout?: ReturnType<typeof setTimeout>;
private schemaValidationTimeout?: ReturnType<typeof setTimeout>;
+ private importPollingSubscription?: Subscription;
readonly selectedFile = signal<File | undefined>(undefined);
readonly uploadId = signal<string | undefined>(undefined);
@@ -117,6 +122,9 @@ export class CsvImportDialogComponent {
CsvImportSchemaValidationResult | undefined
>(undefined);
readonly importResult = signal<CsvImportResult | undefined>(undefined);
+ readonly importJobId = signal<string | undefined>(undefined);
+ readonly importProcessedRows = signal(0);
+ readonly importTotalRows = signal(0);
readonly columnModels = signal<CsvImportColumnModel[]>([]);
readonly previewLoading = signal(false);
readonly importLoading = signal(false);
@@ -274,6 +282,10 @@ export class CsvImportDialogComponent {
this.schedulePreviewReload();
}
});
+
+ this.destroyRef.onDestroy(() => {
+ this.stopImportPolling();
+ });
}
onFileSelected(event: Event): void {
@@ -475,12 +487,14 @@ export class CsvImportDialogComponent {
}
this.importLoading.set(true);
+ this.importProcessedRows.set(0);
+ this.importTotalRows.set(this.previewResult()?.totalRows ?? 0);
this.datalakeRestService
.importCsvData(this.buildImportRequest())
.subscribe({
next: result => {
- this.importResult.set(result);
- this.importLoading.set(false);
+ this.importJobId.set(result.jobId);
+ this.pollImportStatus(result.jobId);
},
error: error => {
this.importLoading.set(false);
@@ -494,7 +508,9 @@ export class CsvImportDialogComponent {
field: 'import',
message:
error?.error?.message ??
- 'CSV import failed.',
+ this.translateService.instant(
+ 'The CSV import failed.',
+ ),
},
],
);
@@ -503,6 +519,7 @@ export class CsvImportDialogComponent {
}
close(refresh = false): void {
+ this.stopImportPolling();
this.dialogRef.close(refresh);
}
@@ -669,10 +686,75 @@ export class CsvImportDialogComponent {
}
private clearImportResult(): void {
+ this.stopImportPolling();
this.importResult.set(undefined);
+ this.importJobId.set(undefined);
+ this.importProcessedRows.set(0);
+ this.importTotalRows.set(0);
this.uploadMessages.set([]);
}
+ private pollImportStatus(jobId: string): void {
+ this.stopImportPolling();
+ this.importPollingSubscription = timer(0, 1000)
+ .pipe(
+ switchMap(() =>
+ this.datalakeRestService.getCsvImportJobStatus(jobId),
+ ),
+ )
+ .subscribe({
+ next: status => this.handleImportStatus(status),
+ error: error => {
+ this.importLoading.set(false);
+ this.stopImportPolling();
+ this.uploadMessages.set([
+ {
+ field: 'import',
+ message:
+ error?.error?.message ??
+ this.translateService.instant(
+ 'CSV import status could not be loaded.',
+ ),
+ },
+ ]);
+ },
+ });
+ }
+
+ private handleImportStatus(status: CsvImportJobStatus): void {
+ this.importProcessedRows.set(status.processedRows);
+ this.importTotalRows.set(status.totalRows);
+
+ if (status.state === 'RUNNING') {
+ return;
+ }
+
+ this.importLoading.set(false);
+ this.stopImportPolling();
+
+ if (status.state === 'SUCCEEDED' && status.result) {
+ this.importResult.set(status.result);
+ } else {
+ this.uploadMessages.set(
+ status.validationMessages?.length
+ ? status.validationMessages
+ : [
+ {
+ field: 'import',
+ message: this.translateService.instant(
+ 'The CSV import failed.',
+ ),
+ },
+ ],
+ );
+ }
+ }
+
+ private stopImportPolling(): void {
+ this.importPollingSubscription?.unsubscribe();
+ this.importPollingSubscription = undefined;
+ }
+
private schedulePreviewReload(): void {
if (!this.currentTarget() || this.previewLoading()) {
return;
diff --git
a/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-upload-state/csv-import-upload-state.component.html
b/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-upload-state/csv-import-upload-state.component.html
index 5acd93e07f..265d4e79b5 100644
---
a/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-upload-state/csv-import-upload-state.component.html
+++
b/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-upload-state/csv-import-upload-state.component.html
@@ -18,16 +18,15 @@
<div class="upload-state">
@if (importLoading()) {
- <mat-spinner [diameter]="56"></mat-spinner>
- <div class="upload-title">
- {{ 'Uploading CSV data' | translate }}
- </div>
- <div class="upload-copy">
- {{
- 'The data is currently being written to the data lake.'
- | translate
- }}
- </div>
+ <sp-progress-bar
+ [title]="'Uploading CSV data' | translate"
+ [ariaLabel]="'CSV import progress' | translate"
+ [value]="processedRows()"
+ [max]="totalRows()"
+ [itemLabel]="'rows imported' | translate"
+ progressBarDataCy="csv-import-progress-bar"
+ progressLabelDataCy="csv-import-progress-label"
+ ></sp-progress-bar>
} @else if (hasImportResult()) {
<mat-icon class="success-icon">check_circle</mat-icon>
<div class="upload-title" data-cy="csv-import-success-title">
diff --git
a/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-upload-state/csv-import-upload-state.component.scss
b/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-upload-state/csv-import-upload-state.component.scss
index 3b740c3d52..cb7507b109 100644
---
a/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-upload-state/csv-import-upload-state.component.scss
+++
b/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-upload-state/csv-import-upload-state.component.scss
@@ -17,31 +17,30 @@
*/
.upload-state {
- min-height: 280px;
+ min-height: 17.5rem;
display: flex;
flex-direction: column;
align-items: center;
justify-content: center;
- gap: 14px;
+ gap: var(--space-md);
text-align: center;
}
.upload-title {
- font-size: 20px;
+ font-size: var(--font-size-xl);
font-weight: 600;
}
-.upload-copy,
.result-meta {
color: var(--color-text-200, #5a6673);
- max-width: 420px;
+ max-width: 26.25rem;
}
.success-icon {
color: var(--color-success, #2e7d32);
- font-size: 48px;
- width: 48px;
- height: 48px;
+ font-size: 3rem;
+ width: 3rem;
+ height: 3rem;
}
.result-name {
@@ -49,10 +48,10 @@
}
.upload-errors {
- max-width: 520px;
+ max-width: 32.5rem;
color: var(--color-text-200, #5a6673);
}
.upload-error-line + .upload-error-line {
- margin-top: 6px;
+ margin-top: var(--space-sm);
}
diff --git
a/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-upload-state/csv-import-upload-state.component.ts
b/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-upload-state/csv-import-upload-state.component.ts
index f2881286a2..df0cc74809 100644
---
a/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-upload-state/csv-import-upload-state.component.ts
+++
b/ui/src/app/dataset/dialog/csv-import-dialog/csv-import-upload-state/csv-import-upload-state.component.ts
@@ -18,13 +18,15 @@
import { Component, input } from '@angular/core';
import { MatIcon } from '@angular/material/icon';
-import { MatProgressSpinner } from '@angular/material/progress-spinner';
import { TranslatePipe } from '@ngx-translate/core';
import {
CsvImportResult,
CsvImportValidationMessage,
} from '@streampipes/platform-services';
-import { SpAlertBannerComponent } from '@streampipes/shared-ui';
+import {
+ ProgressBarComponent,
+ SpAlertBannerComponent,
+} from '@streampipes/shared-ui';
@Component({
selector: 'sp-csv-import-upload-state',
@@ -32,7 +34,7 @@ import { SpAlertBannerComponent } from
'@streampipes/shared-ui';
styleUrls: ['./csv-import-upload-state.component.scss'],
imports: [
MatIcon,
- MatProgressSpinner,
+ ProgressBarComponent,
TranslatePipe,
SpAlertBannerComponent,
],
@@ -42,4 +44,6 @@ export class CsvImportUploadStateComponent {
readonly hasImportResult = input(false);
readonly importResult = input<CsvImportResult | undefined>(undefined);
readonly uploadErrors = input<CsvImportValidationMessage[]>([]);
+ readonly processedRows = input(0);
+ readonly totalRows = input(0);
}