diff --git a/.gitignore b/.gitignore
index 8bb0bd24..0a5f1aa8 100644
--- a/.gitignore
+++ b/.gitignore
@@ -2,7 +2,7 @@
target/classes/*
target/test-classes/*
target/*
-src/main/resources/lib/*.jar
+src/main/resources/jars/*.jar
siesta-query-processor.iml
experiments/*
diff --git a/Dockerfile b/Dockerfile
index a91766f4..d1679aca 100644
--- a/Dockerfile
+++ b/Dockerfile
@@ -2,11 +2,11 @@ FROM ubuntu:20.04
#ENV JAVA_HOME="/usr/lib/jvm/default-jvm/"
-RUN apt-get update && apt-get install -y openjdk-17-jdk maven && \
+RUN apt-get update && apt-get install -y openjdk-17-jdk maven wget && \
echo "export JAVA_HOME=$(dirname $(dirname $(readlink -f $(which java))))" >> /etc/profile.d/java.sh
ENV JAVA_HOME=/usr/lib/jvm/java-17-openjdk-amd64
ENV PATH=$PATH:${JAVA_HOME}/bin
-
+ENV JARS_DIR=/code/src/main/resources/jars
# Install maven
@@ -20,8 +20,15 @@ RUN mvn dependency:resolve
# Adding source, compile and package into a fat jar
ADD src /code/src
+RUN test -f ${JARS_DIR}/hadoop-aws-3.3.4.jar || wget -P ${JARS_DIR} https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-aws/3.3.4/hadoop-aws-3.3.4.jar
+RUN test -f ${JARS_DIR}/aws-java-sdk-bundle-1.12.262.jar || wget -P ${JARS_DIR} https://repo1.maven.org/maven2/com/amazonaws/aws-java-sdk-bundle/1.12.262/aws-java-sdk-bundle-1.12.262.jar
+RUN test -f ${JARS_DIR}/hadoop-client-3.3.4.jar || wget -P ${JARS_DIR} https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-client/3.3.4/hadoop-client-3.3.4.jar
+RUN test -f ${JARS_DIR}/delta-spark_2.12-3.3.0.jar || wget -P ${JARS_DIR} https://repo1.maven.org/maven2/io/delta/delta-spark_2.12/3.3.0/delta-spark_2.12-3.3.0.jar
+RUN test -f ${JARS_DIR}/delta-storage-3.3.0.jar || wget -P ${JARS_DIR} https://repo1.maven.org/maven2/io/delta/delta-storage/3.3.0/delta-storage-3.3.0.jar
RUN mvn clean compile package -f pom.xml -DskipTests
+# Making sure jars are where they should be
+
CMD ["java", "--add-exports", "java.base/sun.nio.ch=ALL-UNNAMED" , "-jar", "target/siesta-query-processor-3.0.jar"]
#ENTRYPOINT ["tail", "-f", "/dev/null"]
diff --git a/README.md b/README.md
index 546d63df..93e75041 100644
--- a/README.md
+++ b/README.md
@@ -41,37 +41,37 @@ To run it locally specify the properties in the application.properties file in t
and then run the project after adding in the ``Add VM options`` under ``Run/Edit configurations`` the following line \
``--add-opens=java.base/sun.nio.ch=ALL-UNNAMED ``
-### Running in docker
-To run it locally open a terminal inside SequenceDetectionQueryExecutor file and run:
-```bash
-docker-compose build
-docker-compose up -d
-```
+### Running with Docker
+
+#### Single Machine Deployment (Development/Testing)
+
+To run everything on one machine:
+
Ensure that this docker and the database can communicate, either run database on a public ip
or connect these two on the same network. You need to specify these environment variables
(you can keep the default ones if you want) in
the docker-compose file before running the QueryExecutor.
+
+#### Distributed Cluster Deployment (Production)
+
+For distributed processing with workers on different machines, see:
+- **[QUICKSTART.md](QUICKSTART.md)** - Quick setup guide
+- **[CLUSTER_DEPLOYMENT.md](CLUSTER_DEPLOYMENT.md)** - Comprehensive deployment documentation
+
+Quick setup:
+
+**On Master Machine:**
+```bash
+./setup-master.sh
```
- master.uri: local[*]
- database: s3
- delta: false # True for streaming, False for batching
- #for s3 (minio)
- s3.endpoint: http://minio:9000
- s3.user: minioadmin
- s3.key: minioadmin
- s3.timeout: 600000
- server.port: 8090
-```
-### SIESTA Query type list
-Below there is a list of all the possible SIESTA queries along with an example JSON or an example url assuming the
-Query Processor is running on localhost:8090
-* GET /health/check (Checks if the application is up and running)
-* GET /lognames (Returns the names of the different log databases)
-* POST /eventTypes (Returns the names of the different event types for a specific log database) \
- Example JSON:
+
+**On Worker Machines:**
+```bash
+./setup-worker.sh
```
-{
- "log_name" : "test"
+
+This will deploy a true distributed Spark cluster with workers running on separate physical machines for better scalability and performance.
+
}
```
* GET /refreshData (Reloads metadata, this should run after a new log file is appended)
diff --git a/docker-compose-swarm.yml b/docker-compose-swarm.yml
new file mode 100644
index 00000000..9922e127
--- /dev/null
+++ b/docker-compose-swarm.yml
@@ -0,0 +1,154 @@
+version: '3'
+
+networks:
+ siesta-swarm-net:
+
+volumes:
+ maven-cache:
+ minio-storage:
+
+services:
+ minio:
+ image: minio/minio:latest
+ container_name: minio
+ environment:
+ MINIO_ROOT_USER: minioadmin
+ MINIO_ROOT_PASSWORD: minioadmin
+ command: server /data
+ ports:
+ - "9000:9000"
+ - "9001:9001"
+ volumes:
+ - minio_storage:/data
+ networks:
+ - siesta-swarm-net
+ deploy:
+ replicas: 1
+ placement:
+ constraints:
+ - 'node.labels.host == m1'
+
+ scylla:
+ container_name: scylla
+ image: scylladb/scylla:latest
+ ports:
+ - "9042:9042"
+ command: --smp 8 --memory 32G --reserve-memory 2G --overprovisioned 1 --api-address 0.0.0.0
+ volumes:
+ - scylla_data:/var/lib/scylla
+ networks:
+ - siesta-net
+ healthcheck:
+ test: ["CMD-SHELL", "nodetool status"]
+ interval: 15s
+ timeout: 15s
+ retries: 5
+ deploy:
+ replicas: 1
+ restart_policy:
+ condition: on-failure
+ max_attempts: 3
+ placement:
+ constraints:
+ - 'node.labels.host == m1'
+
+ query:
+ build: .
+ image: siesta-query:latest
+ stdin_open: true
+ networks:
+ - siesta-swarm-net
+ ports:
+ - '8090:8090'
+ environment:
+ master.uri: spark://spark:7077
+ server.port: 8090
+ database: s3 # cassandra or s3
+ delta: "false" # True for streaming, False for batching
+ spring.mvc.pathmatch.matching-strategy: ANT_PATH_MATCHER
+ #for s3 (minio)
+ s3.endpoint: http://minio:9000
+ s3.user: minioadmin
+ s3.key: minioadmin
+ s3.timeout: 600000
+ # Scylla/Cassandra configuration
+ cassandra.contact.points: scylla
+ cassandra.port: 9042
+ cassandra.keyspace: siesta
+ volumes:
+ - maven-cache:/root/.m2
+ deploy:
+ replicas: 1
+ placement:
+ constraints:
+ - 'node.labels.host == m1'
+
+ spark:
+ image: spark-base:3.5.4
+ networks:
+ - siesta-swarm-net
+ environment:
+ - SPARK_MODE=master
+ - SPARK_MASTER_MEMORY=16G
+ - SPARK_MASTER_CORES=4
+ - SPARK_RPC_AUTHENTICATION_ENABLED=no
+ - SPARK_RPC_ENCRYPTION_ENABLED=no
+ - SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no
+ - SPARK_SSL_ENABLED=no
+ - SPARK_USER=spark
+ - LOGS_PATH=/tmp/spark-events
+ - RESULTS_PATH=/tmp/output/
+ - SPARK_MASTER_WEBUI_PORT=8085
+ ports:
+ - '8085:8085'
+ - '4040:4040'
+ - '8081:8081'
+ deploy:
+ replicas: 1
+ placement:
+ constraints:
+ - 'node.labels.host == m1'
+
+ spark-worker-1:
+ image: spark-base:3.5.4
+ networks:
+ - siesta-swarm-net
+ environment:
+ - SPARK_MODE=worker
+ - SPARK_MASTER_URL=spark://spark:7077
+ - SPARK_WORKER_MEMORY=8G
+ - SPARK_WORKER_CORES=4
+ - SPARK_RPC_AUTHENTICATION_ENABLED=no
+ - SPARK_RPC_ENCRYPTION_ENABLED=no
+ - SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no
+ - SPARK_SSL_ENABLED=no
+ - SPARK_USER=spark
+ #- LOGS_PATH=/tmp/spark-events
+ #- RESULTS_PATH=/tmp/output/
+ deploy:
+ replicas: 1
+ placement:
+ constraints:
+ - 'node.labels.host == m2'
+
+ spark-worker-anaconda-2:
+ image: spark-base:3.5.4
+ networks:
+ - siesta-swarm-net
+ environment:
+ - SPARK_MODE=worker
+ - SPARK_MASTER_URL=spark://spark:7077
+ - SPARK_WORKER_MEMORY=8G
+ - SPARK_WORKER_CORES=4
+ - SPARK_RPC_AUTHENTICATION_ENABLED=no
+ - SPARK_RPC_ENCRYPTION_ENABLED=no
+ - SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no
+ - SPARK_SSL_ENABLED=no
+ - SPARK_USER=spark
+ #- LOGS_PATH=/tmp/spark-events
+ #- RESULTS_PATH=/tmp/output/
+ deploy:
+ replicas: 1
+ placement:
+ constraints:
+ - 'node.labels.host == m2'
diff --git a/docker-compose.yml b/docker-compose.yml
index ac1a1ce9..8eca92b9 100644
--- a/docker-compose.yml
+++ b/docker-compose.yml
@@ -1,78 +1,116 @@
-version: '3.7'
+version: '3.8'
networks:
- cluster_net:
- external: true
+ siesta-net:
name: siesta-net
+ external: true
+ cluster-net:
+ name: cluster-net
+ driver: bridge
+volumes:
+ minio_storage: {}
services:
- funnel:
+ minio:
+ image: minio/minio:latest
+ container_name: minio
+ environment:
+ MINIO_ROOT_USER: minioadmin
+ MINIO_ROOT_PASSWORD: minioadmin
+ command: server /data
+ ports:
+ - "9000:9000"
+ - "9001:9001"
+ volumes:
+ - minio_storage:/data
+ networks:
+ - siesta-net
+
+ scylla:
+ container_name: scylla
+ image: scylladb/scylla:latest
+ ports:
+ - "9042:9042"
+ command: --smp 8 --memory 32G --reserve-memory 2G --overprovisioned 1 --api-address 0.0.0.0
+ volumes:
+ - scylla_data:/var/lib/scylla
+ networks:
+ - siesta-net
+ healthcheck:
+ test: ["CMD-SHELL", "nodetool status"]
+ interval: 15s
+ timeout: 15s
+ retries: 5
+
+ query:
build: .
+ container_name: query
stdin_open: true
environment:
- master.uri: local[*]
- database: s3 # cassandra-rdd or s3
- delta: false # True for streaming, False for batching
- #for s3 (minio)
+ master.uri: spark://spark-master:7077
+ server.port: 8090 # port of the application
+ database: s3 # cassandra or s3
+ delta: "false" # True for streaming, False for batching
+ spring.mvc.pathmatch.matching-strategy: ANT_PATH_MATCHER
+ # S3 (minio)
s3.endpoint: http://minio:9000
s3.user: minioadmin
s3.key: minioadmin
s3.timeout: 600000
- #for cassandra
- cassandra.max_requests_per_local_connection: 32768
- cassandra.max_requests_per_remote_connection: 22000
- cassandra.connections_per_host: 1000
- cassandra.max_queue_size: 1024
- cassandra.connection_timeout: 30000
- cassandra.read_timeout: 30000
- spring.data.cassandra.contact-points: cassandra
- spring.data.cassandra.port: 9042
- spring.data.cassandra.user: cassandra
- spring.data.cassandra.password: cassandra
- server.port: 8090 # port of the application
+ # Scylla/Cassandra
+ cassandra.contact.points: scylla
+ cassandra.port: 9042
+ cassandra.keyspace: siesta
volumes:
- ./build:/root/.m2
ports:
- - '8090:8090'
+ - '8090:8090' # Application port
+ networks:
+ - siesta-net
+ - cluster-net
+
+ spark-master:
+ build:
+ context: .
+ dockerfile: spark-base.Dockerfile
+ image: spark-base:3.5.4
+ container_name: spark-master
+ hostname: spark-master
+ environment:
+ - SPARK_MODE=master
+ - SPARK_MASTER_HOST=0.0.0.0
+ - SPARK_MASTER_PORT=7077
+ - SPARK_MASTER_WEBUI_PORT=8080
+ - SPARK_RPC_AUTHENTICATION_ENABLED=no
+ - SPARK_RPC_ENCRYPTION_ENABLED=no
+ ports:
+ - "7077:7077"
+ - "8080:8080"
networks:
- - cluster_net
+ - cluster-net
+ - siesta-net
-# spark-master:
-# image: bitnami/spark:3.5.4
-# container_name: spark-master
-# environment:
-# - SPARK_MODE=master
-# - SPARK_RPC_AUTHENTICATION_ENABLED=no
-# - SPARK_RPC_ENCRYPTION_ENABLED=no
-# - SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no
-# - SPARK_SSL_ENABLED=no
-# ports:
-# - "7077:7077" # Spark master port
-# - "8080:8080" # Spark Web UI
-# networks:
-# - cluster_net
-#
-# spark-worker-1:
-# image: bitnami/spark:3.5.4
-# container_name: spark-worker-1
-# environment:
-# - SPARK_MODE=worker
-# - SPARK_MASTER_URL=spark://spark-master:7077
-# - SPARK_WORKER_MEMORY=1G
-# - SPARK_WORKER_CORES=1
-# depends_on:
-# - spark-master
-# networks:
-# - cluster_net
-#
-# spark-worker-2:
-# image: bitnami/spark:3.5.4
-# container_name: spark-worker-2
-# environment:
-# - SPARK_MODE=worker
-# - SPARK_MASTER_URL=spark://spark-master:7077
-# - SPARK_WORKER_MEMORY=1G
-# - SPARK_WORKER_CORES=1
-# depends_on:
-# - spark-master
-# networks:
-# - cluster_net
+ spark-worker-1:
+ image: spark-base:3.5.4
+ container_name: spark-worker-1
+ environment:
+ - SPARK_MODE=worker
+ - SPARK_MASTER_URL=spark://spark-master:7077
+ - SPARK_WORKER_MEMORY=2G
+ - SPARK_WORKER_CORES=2
+ depends_on:
+ - spark-master
+ networks:
+ - cluster-net
+
+ spark-worker-2:
+ image: spark-base:3.5.4
+ container_name: spark-worker-2
+ environment:
+ - SPARK_MODE=worker
+ - SPARK_MASTER_URL=spark://spark-master:7077
+ - SPARK_WORKER_MEMORY=2G
+ - SPARK_WORKER_CORES=2
+ depends_on:
+ - spark-master
+ networks:
+ - cluster-net
diff --git a/generate_queries.py b/generate_queries.py
new file mode 100644
index 00000000..eef45356
--- /dev/null
+++ b/generate_queries.py
@@ -0,0 +1,433 @@
+#!/usr/bin/env python3
+# -*- coding: utf-8 -*-
+"""
+Query Generation and Transformation Pipeline
+
+This script provides a complete pipeline for generating queries from XES event logs
+and transforming them into JSON format with random symbol assignments.
+
+Features:
+ - Generate query patterns from XES log files
+ - Extract sequences of specified lengths from log traces
+ - Transform queries to JSON format with symbol assignments
+ - Configurable query lengths and sample sizes
+
+Usage:
+ # Generate queries from XES files and convert to JSON
+ python3 generate_queries.py --mode both --input-dir input --lengths 3,10 --samples 100
+
+ # Only generate .q query files
+ python3 generate_queries.py --mode generate --input-dir input --lengths 3,10 --samples 100
+
+ # Only convert existing .q files to JSON
+ python3 generate_queries.py --mode transform --query-dir queries --output-dir queries_json
+
+Arguments:
+ --mode: Operation mode ('generate', 'transform', or 'both')
+ --input-dir: Directory containing .xes files (default: 'input')
+ --query-dir: Directory for .q query files (default: 'queries')
+ --output-dir: Directory for .jsonl output files (default: 'queries_json')
+ --lengths: Comma-separated list of query lengths (default: '3,10')
+ --samples: Number of samples per length (default: 100)
+
+Output Formats:
+ .q files: LENGTH,EVENT1,EVENT2,...
+ .jsonl files: {"log_name": "...", "pattern": {"eventsWithSymbols": [...]}}
+"""
+
+import argparse
+import json
+import os
+import random
+from typing import List, Tuple
+from statistics import mean, stdev
+
+try:
+ from pm4py.objects.log.importer.xes import importer as xes_import_factory
+ PM4PY_AVAILABLE = True
+except ImportError:
+ PM4PY_AVAILABLE = False
+ print("Warning: pm4py not available. Query generation from XES files will not work.")
+
+
+# Symbol options for query transformation
+SYMBOLS = ["_", "*", "||", "+"]
+
+
+# ============================================================================
+# QUERY GENERATION FROM XES FILES
+# ============================================================================
+
+def generate_query_file(lengths: List[int], samples_per_length: int, log, log_filename: str, output_dir: str):
+ """
+ Generate query file from an XES event log.
+
+ For each specified length, extracts 'samples_per_length' random sequences
+ from traces in the log that are at least that long.
+
+ Args:
+ lengths: List of sequence lengths to generate
+ samples_per_length: Number of query samples to generate per length
+ log: Parsed XES log object from pm4py
+ log_filename: Name of the source log file (used for output naming)
+ output_dir: Directory where .q files will be written
+
+ Returns:
+ Path to the generated query file
+ """
+ os.makedirs(output_dir, exist_ok=True)
+ query_file_path = os.path.join(output_dir, log_filename + ".q")
+
+ total_queries = 0
+ with open(query_file_path, "w", encoding="utf-8") as file:
+ for length in lengths:
+ queries_generated = 0
+ attempts = 0
+ max_attempts = samples_per_length * 100 # Prevent infinite loops
+
+ while queries_generated < samples_per_length and attempts < max_attempts:
+ attempts += 1
+ # Pick a random trace
+ trace_idx = random.randint(0, len(log) - 1)
+ trace = log[trace_idx]
+
+ # Check if trace is long enough
+ if len(trace) >= length:
+ # Extract event names for the first 'length' events
+ event_names = [event["concept:name"] for event in trace][:length]
+ # Write to file: LENGTH,EVENT1,EVENT2,...
+ file.write(str(length) + "," + ",".join(event_names) + "\n")
+ queries_generated += 1
+ total_queries += 1
+
+ if queries_generated < samples_per_length:
+ print(f" Warning: Only generated {queries_generated}/{samples_per_length} queries "
+ f"for length {length} (insufficient long traces)")
+
+ return query_file_path, total_queries
+
+
+def generate_queries_from_xes_files(input_dir: str, query_dir: str, lengths: List[int],
+ samples_per_length: int, force: bool = False):
+ """
+ Process all XES files in input directory and generate query files.
+
+ Args:
+ input_dir: Directory containing .xes files
+ query_dir: Directory where .q query files will be written
+ lengths: List of query lengths to generate
+ samples_per_length: Number of samples per length
+ force: If True, regenerate even if query file already exists
+
+ Returns:
+ Number of query files generated
+ """
+ if not PM4PY_AVAILABLE:
+ print("Error: pm4py is required for query generation. Install with: pip install pm4py")
+ return 0
+
+ if not os.path.isdir(input_dir):
+ print(f"Input directory does not exist: {input_dir}")
+ return 0
+
+ xes_files = [f for f in os.listdir(input_dir) if f.endswith('.xes')]
+ if not xes_files:
+ print(f"No .xes files found in {input_dir}")
+ return 0
+
+ print(f"Found {len(xes_files)} XES file(s) in {input_dir}")
+ print(f"Generating queries with lengths {lengths}, {samples_per_length} samples per length\n")
+
+ files_generated = 0
+ for xes_file in sorted(xes_files):
+ query_file_path = os.path.join(query_dir, xes_file + ".q")
+
+ if os.path.exists(query_file_path) and not force:
+ print(f"Skipping {xes_file} (query file already exists)")
+ continue
+
+ print(f"Processing: {xes_file}")
+ try:
+ # Import XES log
+ log_path = os.path.join(input_dir, xes_file)
+ log = xes_import_factory.apply(log_path)
+ print(f" Loaded log with {len(log)} traces")
+
+ # Generate queries
+ output_path, total_queries = generate_query_file(
+ lengths, samples_per_length, log, xes_file, query_dir
+ )
+ print(f" Generated {total_queries} queries -> {output_path}\n")
+ files_generated += 1
+
+ except Exception as e:
+ print(f" Error processing {xes_file}: {e}\n")
+ continue
+
+ return files_generated
+
+
+# ============================================================================
+# QUERY TRANSFORMATION TO JSON
+# ============================================================================
+
+def transform_query_line(line: str) -> dict:
+ """
+ Parse a single query line and return the pattern object with symbols.
+
+ Input format: LENGTH,EVENT1,EVENT2,...
+ We use the events for ordering/positions and assign random symbols.
+
+ Symbol assignment rules:
+ - Majority of events get "_" (underscore)
+ - At most 2 events get non-underscore symbols ("+", "*", "||")
+ - Ensures semantic variety in query patterns
+
+ Args:
+ line: A line from a .q file
+
+ Returns:
+ Dictionary with "eventsWithSymbols" list, or None if line is invalid
+ """
+ line = line.strip()
+ if not line:
+ return None
+
+ parts = line.split(",")
+ if len(parts) < 2:
+ return None
+
+ # First part is length, rest are event names
+ events = parts[1:]
+
+ # Build initial pattern with all "_" symbols
+ pattern_events = []
+ for idx, event_name in enumerate(events):
+ event_name = event_name.strip()
+ if event_name == "":
+ continue
+ pattern_events.append({
+ "name": event_name,
+ "position": idx,
+ "symbol": "_",
+ })
+
+ if not pattern_events:
+ return None
+
+ # Assign non-underscore symbols to at most 2 events
+ # Ensure majority remain "_"
+ num_events = len(pattern_events)
+ max_non_underscore = min(2, (num_events - 1) // 2)
+
+ if max_non_underscore > 0:
+ # Randomly select 1 to max_non_underscore positions
+ num_to_change = random.randint(1, max_non_underscore)
+ positions_to_change = random.sample(range(num_events), num_to_change)
+ non_underscore_symbols = [s for s in SYMBOLS if s != "_"]
+
+ for pos in positions_to_change:
+ pattern_events[pos]["symbol"] = random.choice(non_underscore_symbols)
+
+ return {"eventsWithSymbols": pattern_events}
+
+
+def transform_query_file(src_path: str, dst_path: str) -> int:
+ """
+ Transform a single .q query file into .jsonl format.
+
+ Args:
+ src_path: Path to input .q file
+ dst_path: Path to output .jsonl file
+
+ Returns:
+ Number of queries written
+ """
+ # Derive log_name from filename
+ basename = os.path.basename(src_path)
+ # Strip .q extension
+ log_name = basename[:-2] if basename.endswith('.q') else basename
+ # Also strip .xes if present (e.g., 'bpi_2017.xes.q' -> 'bpi_2017')
+ if log_name.endswith('.xes'):
+ log_name = log_name[:-4]
+
+ written = 0
+ with open(src_path, "r", encoding="utf-8") as src, \
+ open(dst_path, "w", encoding="utf-8") as dst:
+
+ for line in src:
+ pattern = transform_query_line(line)
+ if pattern is None:
+ continue
+
+ # Create JSON object
+ obj = {
+ "log_name": log_name,
+ "pattern": pattern,
+ }
+ dst.write(json.dumps(obj) + "\n")
+ written += 1
+
+ return written
+
+
+def transform_all_query_files(query_dir: str, output_dir: str) -> Tuple[int, int]:
+ """
+ Transform all .q files in query directory to .jsonl format.
+
+ Args:
+ query_dir: Directory containing .q files
+ output_dir: Directory where .jsonl files will be written
+
+ Returns:
+ Tuple of (files_processed, total_queries_transformed)
+ """
+ if not os.path.isdir(query_dir):
+ print(f"Query directory does not exist: {query_dir}")
+ return 0, 0
+
+ os.makedirs(output_dir, exist_ok=True)
+
+ q_files = [f for f in os.listdir(query_dir) if f.endswith('.q')]
+ if not q_files:
+ print(f"No .q files found in {query_dir}")
+ return 0, 0
+
+ print(f"Found {len(q_files)} .q file(s) in {query_dir}")
+ print(f"Transforming to JSON format with random symbols\n")
+
+ files_processed = 0
+ total_queries = 0
+
+ for q_file in sorted(q_files):
+ src_path = os.path.join(query_dir, q_file)
+ dst_filename = q_file + ".jsonl"
+ dst_path = os.path.join(output_dir, dst_filename)
+
+ try:
+ written = transform_query_file(src_path, dst_path)
+ print(f"Transformed {written} queries: {q_file} -> {dst_filename}")
+ files_processed += 1
+ total_queries += written
+ except Exception as e:
+ print(f"Error processing {q_file}: {e}")
+ continue
+
+ return files_processed, total_queries
+
+
+# ============================================================================
+# MAIN PIPELINE
+# ============================================================================
+
+def main():
+ parser = argparse.ArgumentParser(
+ description="Generate queries from XES logs and transform to JSON format",
+ formatter_class=argparse.RawDescriptionHelpFormatter,
+ epilog="""
+Examples:
+ # Full pipeline: generate and transform
+ python3 generate_queries.py --mode both --input-dir input --lengths 3,10 --samples 100
+
+ # Only generate .q files from XES logs
+ python3 generate_queries.py --mode generate --input-dir input
+
+ # Only transform existing .q files to JSON
+ python3 generate_queries.py --mode transform --query-dir queries
+
+ # Custom configuration
+ python3 generate_queries.py --mode both --lengths 3,5,10,15 --samples 50 --force
+ """
+ )
+
+ parser.add_argument(
+ '--mode',
+ choices=['generate', 'transform', 'both'],
+ default='both',
+ help='Operation mode: generate .q files, transform to JSON, or both (default: both)'
+ )
+ parser.add_argument(
+ '--input-dir',
+ default='input',
+ help='Directory containing .xes files (default: input)'
+ )
+ parser.add_argument(
+ '--query-dir',
+ default='queries',
+ help='Directory for .q query files (default: queries)'
+ )
+ parser.add_argument(
+ '--output-dir',
+ default='queries_json',
+ help='Directory for .jsonl output files (default: queries_json)'
+ )
+ parser.add_argument(
+ '--lengths',
+ default='3,10',
+ help='Comma-separated list of query lengths (default: 3,10)'
+ )
+ parser.add_argument(
+ '--samples',
+ type=int,
+ default=100,
+ help='Number of query samples per length (default: 100)'
+ )
+ parser.add_argument(
+ '--force',
+ action='store_true',
+ help='Force regeneration of existing query files'
+ )
+
+ args = parser.parse_args()
+
+ # Parse lengths
+ try:
+ lengths = [int(x.strip()) for x in args.lengths.split(',')]
+ except ValueError:
+ print(f"Error: Invalid lengths format '{args.lengths}'. Use comma-separated integers.")
+ return 1
+
+ # Convert to absolute paths
+ input_dir = os.path.abspath(args.input_dir)
+ query_dir = os.path.abspath(args.query_dir)
+ output_dir = os.path.abspath(args.output_dir)
+
+ print("=" * 70)
+ print("QUERY GENERATION AND TRANSFORMATION PIPELINE")
+ print("=" * 70)
+ print(f"Mode: {args.mode}")
+ print(f"Input directory: {input_dir}")
+ print(f"Query directory: {query_dir}")
+ print(f"Output directory: {output_dir}")
+ print(f"Query lengths: {lengths}")
+ print(f"Samples per length: {args.samples}")
+ print("=" * 70)
+ print()
+
+ # Execute based on mode
+ if args.mode in ['generate', 'both']:
+ print("STEP 1: Generating queries from XES files")
+ print("-" * 70)
+ files_generated = generate_queries_from_xes_files(
+ input_dir, query_dir, lengths, args.samples, args.force
+ )
+ print(f"Summary: Generated query files from {files_generated} XES log(s)")
+ print()
+
+ if args.mode in ['transform', 'both']:
+ print("STEP 2: Transforming queries to JSON format")
+ print("-" * 70)
+ files_processed, total_queries = transform_all_query_files(query_dir, output_dir)
+ print()
+ print(f"Summary: Processed {files_processed} file(s), transformed {total_queries} queries")
+ print()
+
+ print("=" * 70)
+ print("Pipeline completed successfully!")
+ print("=" * 70)
+
+ return 0
+
+
+if __name__ == '__main__':
+ raise SystemExit(main())
diff --git a/pom.xml b/pom.xml
index ae366779..8292b553 100644
--- a/pom.xml
+++ b/pom.xml
@@ -223,11 +223,36 @@
-
-
-
-
-
+
+ io.delta
+ delta-core_2.12
+ 2.1.0
+
+
+
+
+ com.datastax.spark
+ spark-cassandra-connector_2.12
+ 3.4.1
+
+
+
+ com.datastax.oss
+ java-driver-core
+ 4.17.0
+
+
+
+ com.amazonaws
+ aws-java-sdk-bundle
+ 1.12.262
+
+
+
+ com.datastax.oss
+ java-driver-query-builder
+ 4.17.0
+
io.delta
@@ -239,6 +264,16 @@
+
+ org.apache.maven.plugins
+ maven-compiler-plugin
+ 3.11.0
+
+ 17
+ 17
+ 17
+
+
org.springframework.boot
spring-boot-maven-plugin
diff --git a/spark-base.Dockerfile b/spark-base.Dockerfile
new file mode 100644
index 00000000..68f840f6
--- /dev/null
+++ b/spark-base.Dockerfile
@@ -0,0 +1,13 @@
+FROM public.ecr.aws/bitnami/spark:3.5.4
+
+USER root
+
+# Install curl (Bitnami images are Debian-based)
+RUN apt-get update && apt-get install -y curl && rm -rf /var/lib/apt/lists/*
+
+# Download AWS + Hadoop JARs for S3 support
+RUN mkdir -p /opt/bitnami/spark/jars && \
+ curl -L -o /opt/bitnami/spark/jars/hadoop-aws-3.3.4.jar https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-aws/3.3.4/hadoop-aws-3.3.4.jar && \
+ curl -L -o /opt/bitnami/spark/jars/aws-java-sdk-bundle-1.12.262.jar https://repo1.maven.org/maven2/com/amazonaws/aws-java-sdk-bundle/1.12.262/aws-java-sdk-bundle-1.12.262.jar
+
+USER 1001
diff --git a/src/main/java/com/datalab/siesta/queryprocessor/SaseConnection/SaseConnector.java b/src/main/java/com/datalab/siesta/queryprocessor/SaseConnection/SaseConnector.java
index 9971a781..476d1519 100644
--- a/src/main/java/com/datalab/siesta/queryprocessor/SaseConnection/SaseConnector.java
+++ b/src/main/java/com/datalab/siesta/queryprocessor/SaseConnection/SaseConnector.java
@@ -64,10 +64,18 @@ public List evaluate(SIESTAPattern pattern, Map
Occurrences ocs = new Occurrences();
ocs.setTraceID(e.getKey());
for (Match m : ec.getMatches()) {
- ocs.addOccurrence(new Occurrence(Arrays.stream(m.getEvents()).parallel()
- .map(x -> (SaseEvent) x)
- .map(SaseEvent::getEventBoth)
- .collect(Collectors.toList())));
+ List eventBothList = new ArrayList<>();
+ try {
+ for (edu.umass.cs.sase.stream.Event event : m.getEvents()) {
+ SaseEvent sEvent = (SaseEvent) event;
+ EventBoth both = sEvent.getEventBoth();
+ eventBothList.add(both);
+ }
+ Occurrence oc = new Occurrence(eventBothList);
+ ocs.addOccurrence(oc);
+ }catch (Exception ignored){
+
+ }
}
occurrences.add(ocs);
}
@@ -98,10 +106,18 @@ public List evaluateGroups(SIESTAPattern pattern, Map (SaseEvent) x)
- .map(SaseEvent::getEventBoth)
- .collect(Collectors.toList())));
+ List eventBothList = new ArrayList<>();
+ try {
+ for (edu.umass.cs.sase.stream.Event event : m.getEvents()) {
+ SaseEvent sEvent = (SaseEvent) event;
+ EventBoth both = sEvent.getEventBoth();
+ eventBothList.add(both);
+ }
+ Occurrence oc = new Occurrence(eventBothList);
+ ocs.addOccurrence(oc);
+ }catch (Exception ignored){
+
+ }
}
occurrences.add(ocs);
}
@@ -132,10 +148,18 @@ public List evaluateSmallPatterns(SIESTAPattern pattern, Map (SaseEvent) x)
- .map(SaseEvent::getEventBoth)
- .collect(Collectors.toList())));
+ List eventBothList = new ArrayList<>();
+ try {
+ for (edu.umass.cs.sase.stream.Event event : m.getEvents()) {
+ SaseEvent sEvent = (SaseEvent) event;
+ EventBoth both = sEvent.getEventBoth();
+ eventBothList.add(both);
+ }
+ Occurrence oc = new Occurrence(eventBothList);
+ ocs.addOccurrence(oc);
+ }catch (Exception ignored){
+
+ }
}
occurrences.add(ocs);
}
diff --git a/src/main/java/com/datalab/siesta/queryprocessor/controllers/PatternAnalysisController.java b/src/main/java/com/datalab/siesta/queryprocessor/controllers/PatternAnalysisController.java
index 6742273e..305dbfcb 100644
--- a/src/main/java/com/datalab/siesta/queryprocessor/controllers/PatternAnalysisController.java
+++ b/src/main/java/com/datalab/siesta/queryprocessor/controllers/PatternAnalysisController.java
@@ -176,7 +176,7 @@ private Tuple2, List> splitLogInstances(String log_da
List legitimateInstances = new ArrayList<>();
List illegitimateInstances = new ArrayList<>();
indexRecords.forEach((et, instances) -> instances.forEach(instance -> {
- if (legitimateTraces.contains(instance.getTraceId())) {
+ if (legitimateTraces.contains(instance.getTrace_id())) {
legitimateInstances.add(instance);
} else {
illegitimateInstances.add(instance);
@@ -192,8 +192,8 @@ private Tuple2, List> splitLogInstances(String log_da
*/
private double getCV(Count pair) {
double norm_factor = 1.0 / pair.getCount();
- double mean = (double) pair.getSum_duration() / pair.getCount();
- double var = norm_factor * (pair.getSum_squares() - norm_factor * Math.pow(pair.getSum_duration(), 2));
+ double mean = (double) pair.getSumDuration() / pair.getCount();
+ double var = norm_factor * (pair.getSumSquares() - norm_factor * Math.pow(pair.getSumDuration(), 2));
return var / mean;
}
diff --git a/src/main/java/com/datalab/siesta/queryprocessor/declare/DeclareDBConnector.java b/src/main/java/com/datalab/siesta/queryprocessor/declare/DeclareDBConnector.java
index e61f9834..cd446600 100644
--- a/src/main/java/com/datalab/siesta/queryprocessor/declare/DeclareDBConnector.java
+++ b/src/main/java/com/datalab/siesta/queryprocessor/declare/DeclareDBConnector.java
@@ -1,31 +1,22 @@
package com.datalab.siesta.queryprocessor.declare;
-import com.datalab.siesta.queryprocessor.declare.model.EventPairToTrace;
-import com.datalab.siesta.queryprocessor.declare.model.EventSupport;
-import com.datalab.siesta.queryprocessor.declare.model.OccurrencesPerTrace;
-import com.datalab.siesta.queryprocessor.declare.model.UniqueTracesPerEventPair;
-import com.datalab.siesta.queryprocessor.declare.model.UniqueTracesPerEventType;
+import com.datalab.siesta.queryprocessor.declare.model.*;
import com.datalab.siesta.queryprocessor.declare.model.declareState.ExistenceState;
import com.datalab.siesta.queryprocessor.declare.model.declareState.NegativeState;
import com.datalab.siesta.queryprocessor.declare.model.declareState.OrderState;
import com.datalab.siesta.queryprocessor.declare.model.declareState.PositionState;
import com.datalab.siesta.queryprocessor.declare.model.declareState.UnorderStateI;
import com.datalab.siesta.queryprocessor.declare.model.declareState.UnorderStateU;
-import com.datalab.siesta.queryprocessor.model.DBModel.IndexPair;
-import com.datalab.siesta.queryprocessor.model.DBModel.Trace;
-import com.datalab.siesta.queryprocessor.model.Events.EventBoth;
-import com.datalab.siesta.queryprocessor.model.Events.EventPos;
+import com.datalab.siesta.queryprocessor.storage.model.EventTypeTracePositions;
+import com.datalab.siesta.queryprocessor.storage.model.Trace;
import com.datalab.siesta.queryprocessor.storage.DatabaseRepository;
-import org.apache.spark.api.java.JavaPairRDD;
-import org.apache.spark.api.java.JavaRDD;
+
+import org.apache.spark.sql.Dataset;
+import org.apache.spark.sql.Encoders;
+import org.apache.spark.sql.functions;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
-import scala.Tuple2;
-import scala.Tuple3;
-
-import java.util.List;
-import java.util.Map;
@Service
public class DeclareDBConnector {
@@ -37,62 +28,60 @@ public DeclareDBConnector(DatabaseRepository databaseRepository){
this.db=databaseRepository;
}
- public JavaRDD querySequenceTableDeclare(String logName){
+ public Dataset querySequenceTableDeclare(String logName){
return db.querySequenceTableDeclare(logName);
}
- public JavaRDD querySingleTableDeclare(String logname){
+ public Dataset querySingleTableDeclare(String logname){
return this.db.querySingleTableDeclare(logname);
}
- public JavaRDD querySingleTable(String logname){
+ public Dataset querySingleTable(String logname){
return this.db.querySingleTable(logname);
}
- public JavaRDD queryIndexTableDeclare(String logname){
+ public Dataset queryIndexTableDeclare(String logname){
return this.db.queryIndexTableDeclare(logname);
}
- public JavaRDD queryIndexTableAllDeclare(String logname){
- return this.db.queryIndexTableAllDeclare(logname);
- }
-
- public JavaPairRDD, List> querySingleTableAllDeclare(String logname){
+ public Dataset querySingleTableAllDeclare(String logname){
return this.db.querySingleTableAllDeclare(logname);
}
- public JavaRDD queryIndexOriginalDeclare(String logname){
+ public Dataset queryIndexOriginalDeclare(String logname){
return this.db.queryIndexOriginalDeclare(logname);
}
- public Map extractTotalOccurrencesPerEventType(String logname){
- return this.querySingleTableDeclare(logname)
- .map(x -> {
- long all = x.getOccurrences().stream().mapToLong(OccurrencesPerTrace::getOccurrences).sum();
- return new Tuple2<>(x.getEventType(), all);
- }).keyBy(x -> x._1).mapValues(x -> x._2).collectAsMap();
+ public Dataset extractTotalOccurrencesPerEventType(String logname){
+ Dataset eventTypeOccurrencesDataset= this.querySingleTableDeclare(logname)
+ .withColumn("numberOfTraces", functions.expr(
+ "aggregate(occurrences, 0, (acc, a) -> acc + a.occs)"
+ ))
+ .selectExpr("eventName", "numberOfTraces")
+ .as(Encoders.bean(EventTypeOccurrences.class));
+ return eventTypeOccurrencesDataset;
}
- public JavaRDD queryPositionState(String logname){
+ public Dataset queryPositionState(String logname){
return this.db.queryPositionState(logname);
}
- public JavaRDD queryExistenceState(String logname){
+ public Dataset queryExistenceState(String logname){
return this.db.queryExistenceState(logname);
}
- public JavaRDD queryUnorderStateI(String logname){
+ public Dataset queryUnorderStateI(String logname){
return this.db.queryUnorderStateI(logname);
}
- public JavaRDD queryUnorderStateU(String logname){
+ public Dataset queryUnorderStateU(String logname){
return this.db.queryUnorderStateU(logname);
}
- public JavaRDD queryOrderState(String logname){
+ public Dataset queryOrderState(String logname){
return this.db.queryOrderState(logname);
}
- public JavaRDD queryNegativeState(String logname){
+ public Dataset queryNegativeState(String logname){
return this.db.queryNegativeState(logname);
}
diff --git a/src/main/java/com/datalab/siesta/queryprocessor/declare/DeclareUtilities.java b/src/main/java/com/datalab/siesta/queryprocessor/declare/DeclareUtilities.java
index 5d198108..645e7e71 100644
--- a/src/main/java/com/datalab/siesta/queryprocessor/declare/DeclareUtilities.java
+++ b/src/main/java/com/datalab/siesta/queryprocessor/declare/DeclareUtilities.java
@@ -1,9 +1,9 @@
package com.datalab.siesta.queryprocessor.declare;
import com.datalab.siesta.queryprocessor.declare.model.EventPairToNumberOfTrace;
-import com.datalab.siesta.queryprocessor.model.Events.Event;
-import com.datalab.siesta.queryprocessor.model.Events.EventPair;
-import org.apache.spark.api.java.JavaRDD;
+import com.datalab.siesta.queryprocessor.model.DBModel.EventTypes;
+import org.apache.spark.sql.Dataset;
+import org.apache.spark.sql.Encoders;
import org.jvnet.hk2.annotations.Service;
import org.springframework.stereotype.Component;
@@ -22,21 +22,24 @@ public DeclareUtilities() {
* @param joined a rdd containing all the event pairs that occurred
* @return a set of all the event pairs that did not appear in the log database
*/
- public Set extractNotFoundPairs(Set eventTypes,JavaRDD joined) {
+ public Set extractNotFoundPairs(Set eventTypes, Dataset joined) {
//calculate all the event pairs (n^2) and store them in a set
//event pairs of type (eventA,eventA) are excluded
- Set allEventPairs = new HashSet<>();
+ Set allEventPairs = new HashSet<>();
for (String eventA : eventTypes) {
for (String eventB : eventTypes) {
if (!eventA.equals(eventB)) {
- allEventPairs.add(new EventPair(new Event(eventA), new Event(eventB)));
+ allEventPairs.add(new EventTypes(eventA, eventB));
}
}
}
//removes from the above set all the event pairs that have at least one occurrence in the
- List foundEventPairs = joined.map(x -> new EventPair(new Event(x.getEventA()), new Event(x.getEventB())))
- .collect();
+ List foundEventPairs = joined
+ .select("eventA","eventB")
+ .as(Encoders.bean(EventTypes.class))
+ .collectAsList();
+
foundEventPairs.forEach(allEventPairs::remove);
return allEventPairs;
}
diff --git a/src/main/java/com/datalab/siesta/queryprocessor/declare/model/EventSupport.java b/src/main/java/com/datalab/siesta/queryprocessor/declare/model/EventSupport.java
index b0398ac9..a75d929b 100644
--- a/src/main/java/com/datalab/siesta/queryprocessor/declare/model/EventSupport.java
+++ b/src/main/java/com/datalab/siesta/queryprocessor/declare/model/EventSupport.java
@@ -3,7 +3,15 @@
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.databind.annotation.JsonSerialize;
-
+import lombok.AllArgsConstructor;
+import lombok.Getter;
+import lombok.NoArgsConstructor;
+import lombok.Setter;
+
+@Getter
+@Setter
+@AllArgsConstructor
+@NoArgsConstructor
public class EventSupport {
@JsonProperty("ev")
@@ -11,28 +19,5 @@ public class EventSupport {
@JsonProperty("support")
@JsonSerialize(using = SupportSerializer.class)
protected Double support;
-
- public EventSupport(String event, double support) {
- this.event = event;
- this.support = support;
- }
-
- public EventSupport() {
- }
-
- public String getEvent() {
- return event;
- }
-
- public void setEvent(String event) {
- this.event = event;
- }
-
- public double getSupport() {
- return support;
- }
-
- public void setSupport(double support) {
- this.support = support;
- }
}
+
diff --git a/src/main/java/com/datalab/siesta/queryprocessor/declare/model/EventTypeOccurrences.java b/src/main/java/com/datalab/siesta/queryprocessor/declare/model/EventTypeOccurrences.java
new file mode 100644
index 00000000..be8a52ea
--- /dev/null
+++ b/src/main/java/com/datalab/siesta/queryprocessor/declare/model/EventTypeOccurrences.java
@@ -0,0 +1,17 @@
+package com.datalab.siesta.queryprocessor.declare.model;
+
+import lombok.AllArgsConstructor;
+import lombok.Getter;
+import lombok.NoArgsConstructor;
+import lombok.Setter;
+
+import java.io.Serializable;
+
+@Getter
+@Setter
+@NoArgsConstructor
+@AllArgsConstructor
+public class EventTypeOccurrences implements Serializable {
+ private String eventName;
+ private long numberOfTraces;
+}
diff --git a/src/main/java/com/datalab/siesta/queryprocessor/declare/model/OccurrencesPerTrace.java b/src/main/java/com/datalab/siesta/queryprocessor/declare/model/OccurrencesPerTrace.java
index e8f9adac..437f28f7 100644
--- a/src/main/java/com/datalab/siesta/queryprocessor/declare/model/OccurrencesPerTrace.java
+++ b/src/main/java/com/datalab/siesta/queryprocessor/declare/model/OccurrencesPerTrace.java
@@ -17,5 +17,5 @@
@NoArgsConstructor
public class OccurrencesPerTrace implements Serializable {
private String traceId;
- private int occurrences;
+ private int occs;
}
diff --git a/src/main/java/com/datalab/siesta/queryprocessor/declare/model/UniqueTracesPerEventPair.java b/src/main/java/com/datalab/siesta/queryprocessor/declare/model/UniqueTracesPerEventPair.java
index 5aa70b6a..e3254033 100644
--- a/src/main/java/com/datalab/siesta/queryprocessor/declare/model/UniqueTracesPerEventPair.java
+++ b/src/main/java/com/datalab/siesta/queryprocessor/declare/model/UniqueTracesPerEventPair.java
@@ -1,47 +1,24 @@
package com.datalab.siesta.queryprocessor.declare.model;
import com.datalab.siesta.queryprocessor.model.Events.EventPair;
+import lombok.AllArgsConstructor;
+import lombok.Getter;
+import lombok.NoArgsConstructor;
+import lombok.Setter;
import scala.Tuple2;
import java.io.Serializable;
import java.util.List;
-
+@Getter
+@Setter
+@AllArgsConstructor
+@NoArgsConstructor
public class UniqueTracesPerEventPair implements Serializable {
private String eventA;
private String eventB;
private List uniqueTraces;
- public UniqueTracesPerEventPair(String eventA, String eventB, List uniqueTraces) {
- this.eventA = eventA;
- this.eventB = eventB;
- this.uniqueTraces = uniqueTraces;
- }
-
- public String getEventA() {
- return eventA;
- }
-
- public void setEventA(String eventA) {
- this.eventA = eventA;
- }
-
- public String getEventB() {
- return eventB;
- }
-
- public void setEventB(String eventB) {
- this.eventB = eventB;
- }
-
- public List getUniqueTraces() {
- return uniqueTraces;
- }
-
- public void setUniqueTraces(List uniqueTraces) {
- this.uniqueTraces = uniqueTraces;
- }
-
public Tuple2 getKey(){
return new Tuple2<>(this.eventA,this.eventB);
}
diff --git a/src/main/java/com/datalab/siesta/queryprocessor/declare/model/UniqueTracesPerEventType.java b/src/main/java/com/datalab/siesta/queryprocessor/declare/model/UniqueTracesPerEventType.java
index 43c3cae8..e395bc0b 100644
--- a/src/main/java/com/datalab/siesta/queryprocessor/declare/model/UniqueTracesPerEventType.java
+++ b/src/main/java/com/datalab/siesta/queryprocessor/declare/model/UniqueTracesPerEventType.java
@@ -19,7 +19,7 @@
@Getter
@Setter
public class UniqueTracesPerEventType implements Serializable {
- private String eventType;
+ private String eventName;
private List occurrences;
@@ -31,10 +31,10 @@ public class UniqueTracesPerEventType implements Serializable {
public HashMap groupTimes(){
HashMap groups = new HashMap<>();
this.occurrences.forEach(x->{
- if(groups.containsKey(x.getOccurrences())){
- groups.put(x.getOccurrences(),groups.get(x.getOccurrences())+1L);
+ if(groups.containsKey(x.getOccs())){
+ groups.put(x.getOccs(),groups.get(x.getOccs())+1L);
}else{
- groups.put(x.getOccurrences(),1L);
+ groups.put(x.getOccs(),1L);
}
});
return groups;
diff --git a/src/main/java/com/datalab/siesta/queryprocessor/declare/model/declareState/ExistenceState.java b/src/main/java/com/datalab/siesta/queryprocessor/declare/model/declareState/ExistenceState.java
index ae33579b..68699fce 100644
--- a/src/main/java/com/datalab/siesta/queryprocessor/declare/model/declareState/ExistenceState.java
+++ b/src/main/java/com/datalab/siesta/queryprocessor/declare/model/declareState/ExistenceState.java
@@ -3,18 +3,16 @@
import java.io.Serializable;
import lombok.Getter;
+import lombok.NoArgsConstructor;
import lombok.Setter;
@Setter
@Getter
+@NoArgsConstructor
public class ExistenceState implements Serializable{
private String event_type;
private int occurrences;
private long contained;
- public ExistenceState(){
-
- }
-
}
diff --git a/src/main/java/com/datalab/siesta/queryprocessor/declare/queryPlans/QueryPlanDeclareAll.java b/src/main/java/com/datalab/siesta/queryprocessor/declare/queryPlans/QueryPlanDeclareAll.java
index 8703e294..54eb5327 100644
--- a/src/main/java/com/datalab/siesta/queryprocessor/declare/queryPlans/QueryPlanDeclareAll.java
+++ b/src/main/java/com/datalab/siesta/queryprocessor/declare/queryPlans/QueryPlanDeclareAll.java
@@ -18,16 +18,15 @@
import com.datalab.siesta.queryprocessor.model.Queries.QueryResponses.QueryResponse;
import com.datalab.siesta.queryprocessor.model.Queries.Wrapper.QueryWrapper;
+import com.datalab.siesta.queryprocessor.storage.model.EventTypeTracePositions;
import lombok.Setter;
-import org.apache.spark.api.java.JavaPairRDD;
-import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.api.java.JavaSparkContext;
-import org.apache.spark.broadcast.Broadcast;
+import org.apache.spark.sql.Dataset;
+import org.apache.spark.sql.functions;
import org.apache.spark.storage.StorageLevel;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import org.springframework.web.context.annotation.RequestScope;
-import scala.Tuple2;
import java.util.*;
@@ -72,51 +71,47 @@ public QueryResponse execute(QueryWrapper qw) {
// run existences
this.queryPlanExistences.setMetadata(metadata);
this.queryPlanExistences.initResponse();
- Broadcast bSupport = javaSparkContext.broadcast(support);
- Broadcast bTotalTraces = javaSparkContext.broadcast(metadata.getTraces());
// if existence, absence or exactly in modes
- JavaRDD uEventType = declareDBConnector
+ Dataset uEventType = declareDBConnector
.querySingleTableDeclare(this.metadata.getLogname());
uEventType.persist(StorageLevel.MEMORY_AND_DISK());
Map> groupTimes = this.queryPlanExistences.createMapForSingle(uEventType);
Map singleUnique = this.queryPlanExistences.extractUniqueTracesSingle(groupTimes);
- Broadcast