diff --git a/src/main/java/org/rumbledb/api/SequenceWriter.java b/src/main/java/org/rumbledb/api/SequenceWriter.java index 9576e88ae8..560ff8e7d0 100644 --- a/src/main/java/org/rumbledb/api/SequenceWriter.java +++ b/src/main/java/org/rumbledb/api/SequenceWriter.java @@ -4,8 +4,7 @@ import java.util.List; import java.util.Map; -import org.apache.log4j.LogManager; -import org.apache.log4j.Logger; +import lombok.extern.log4j.Log4j2; import org.apache.spark.api.java.JavaRDD; import org.apache.spark.sql.DataFrameWriter; import org.apache.spark.sql.Dataset; @@ -43,6 +42,8 @@ * {@link SerializationParameters#getMethod()}, which is the single source of truth for the output * format. */ + +@Log4j2 public class SequenceWriter { private static final int SINGLE_PARTITION_CAP = 1000000000; @@ -287,11 +288,10 @@ public void save(String path) { // DataFrame mode: delegate to Spark's DataFrameWriter, using the serialization method // as the Spark output format (json/csv/parquet/other). if (this.dataFrameWriter != null) { - Logger logger = LogManager.getLogger(SequenceWriter.class); for (Map.Entry option : this.serializationParameters.getSparkOptions().entrySet()) { - logger.info("Writing with option " + option.getKey() + " : " + option.getValue()); + log.info("Writing with option " + option.getKey() + " : " + option.getValue()); } - logger.info("Writing to format " + method); + log.info("Writing to format " + method); DataFrameWriter writerWithOptions = applyStoredSparkOptions(this.dataFrameWriter); String target = FileSystemUtil.convertURIToStringForSpark(outputUri); if (method.equalsIgnoreCase("json")) { diff --git a/src/main/java/org/rumbledb/cli/LoggingConfiguration.java b/src/main/java/org/rumbledb/cli/LoggingConfiguration.java new file mode 100644 index 0000000000..52c5f9c59d --- /dev/null +++ b/src/main/java/org/rumbledb/cli/LoggingConfiguration.java @@ -0,0 +1,85 @@ +/* + * 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.rumbledb.cli; + +import java.util.Arrays; +import java.util.Locale; + +import org.apache.logging.log4j.Level; +import org.apache.logging.log4j.core.config.Configurator; +import org.apache.logging.log4j.spi.StandardLevel; +import org.rumbledb.config.model.DebugConfig; +import org.rumbledb.exceptions.CliException; +import org.rumbledb.spark.SparkSessionManager; + +final class LoggingConfiguration { + private static final String LOGGER_NAME = "org.rumbledb"; + private static final Level DEFAULT_APPLICATION_LEVEL = Level.WARN; + private static final Level DEFAULT_SPARK_LEVEL = Level.OFF; + private static final String VALID_LEVELS = String.join( + ", ", + Arrays.stream(StandardLevel.values()) + .map(StandardLevel::name) + .map(level -> level.toLowerCase(Locale.ROOT)) + .toList() + ); + + private LoggingConfiguration() { + } + + static void configure(DebugConfig debugConfig) { + Configurator.setLevel(LOGGER_NAME, resolveApplicationLevel(debugConfig)); + SparkSessionManager.setLogLevel(resolveSparkLevel(debugConfig)); + } + + private static Level resolveApplicationLevel(DebugConfig debugConfig) { + if (isBlank(debugConfig.logLevel())) { + return DEFAULT_APPLICATION_LEVEL; + } + return parseLevel(debugConfig.logLevel(), "--log-level"); + } + + private static Level resolveSparkLevel(DebugConfig debugConfig) { + if (isBlank(debugConfig.sparkLogLevel())) { + return DEFAULT_SPARK_LEVEL; + } + return parseLevel(debugConfig.sparkLogLevel(), "--spark-log-level"); + } + + private static Level parseLevel(String rawLevel, String optionName) { + String normalizedLevel = rawLevel.trim().toUpperCase(Locale.ROOT); + try { + StandardLevel.valueOf(normalizedLevel); + } catch (IllegalArgumentException e) { + throw new CliException( + "Invalid " + + optionName + + " value: " + + rawLevel + + ". Valid values are: " + + VALID_LEVELS + + "." + ); + } + return Level.toLevel(normalizedLevel); + } + + private static boolean isBlank(String value) { + return value == null || value.isBlank(); + } +} diff --git a/src/main/java/org/rumbledb/cli/Main.java b/src/main/java/org/rumbledb/cli/Main.java index 4df66181b8..b34647513e 100644 --- a/src/main/java/org/rumbledb/cli/Main.java +++ b/src/main/java/org/rumbledb/cli/Main.java @@ -61,9 +61,10 @@ public static void main(String[] args) throws IOException { System.exit(0); } else { configuration = invocation.configuration(); + LoggingConfiguration.configure(configuration.debug()); } } catch (Exception e) { - System.err.println("⚠️ CLI Error: " + e.getMessage()); + ConsoleOutput.error("⚠️ CLI Error: " + e.getMessage()); System.exit(42); } diff --git a/src/main/java/org/rumbledb/cli/arguments/DebugArguments.java b/src/main/java/org/rumbledb/cli/arguments/DebugArguments.java index 7986de6f0d..641fabab69 100644 --- a/src/main/java/org/rumbledb/cli/arguments/DebugArguments.java +++ b/src/main/java/org/rumbledb/cli/arguments/DebugArguments.java @@ -29,12 +29,28 @@ public final class DebugArguments { ) private Boolean logging; + @Option( + names = "--log-level", + paramLabel = "level", + description = "Sets the diagnostic logging level. Valid values: off, fatal, error, warn, info, debug, trace, all." + ) + private String logLevel; + + @Option( + names = "--spark-log-level", + paramLabel = "level", + description = "Sets the Spark logging level. Valid values: off, fatal, error, warn, info, debug, trace, all." + ) + private String sparkLogLevel; + public DebugConfig toConfig() { DebugConfig.DebugConfigBuilder builder = DebugConfig.builder(); OptionConversion.applyBooleanIfPresent(this.printIteratorTree, builder::printIteratorTree); OptionConversion.applyBooleanIfPresent(this.showErrorInfo, builder::showErrorInfo); OptionConversion.applyBooleanIfPresent(this.logging, builder::logging); + OptionConversion.applyIfPresent(this.logLevel, builder::logLevel); + OptionConversion.applyIfPresent(this.sparkLogLevel, builder::sparkLogLevel); return builder.build(); } diff --git a/src/main/java/org/rumbledb/compiler/ExecutionModeVisitor.java b/src/main/java/org/rumbledb/compiler/ExecutionModeVisitor.java index 8225843d13..891fd1c03b 100644 --- a/src/main/java/org/rumbledb/compiler/ExecutionModeVisitor.java +++ b/src/main/java/org/rumbledb/compiler/ExecutionModeVisitor.java @@ -20,7 +20,7 @@ package org.rumbledb.compiler; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.rumbledb.bindings.DataFrameBinding; import org.rumbledb.bindings.ExternalBindings; import org.rumbledb.config.RumbleConfiguration; @@ -96,6 +96,8 @@ /** * Static context visitor implements a multi-pass algorithm that enables function hoisting */ + +@Log4j2 public class ExecutionModeVisitor extends AbstractNodeVisitor { private VisitorConfig visitorConfig; @@ -162,18 +164,15 @@ public StaticContext visitVariableReference(VariableReferenceExpression expressi expression.getStaticSequenceType().getArity().equals(Arity.ZeroOrMore) ) { if (expression.getStaticSequenceType().getItemType().isObjectItemType()) { - System.err.println( - "[WARNING] Forcing execution mode of variable " - + expression.getVariableName() - + " to DataFrame based on static object* type." + log.warn( + "Forcing execution mode of variable {} to DataFrame based on static object* type.", + expression.getVariableName() ); expression.setHighestExecutionMode(DATAFRAMEifConfigurationAllows()); return argument; } } - System.err.println( - "[WARNING] Forcing execution mode of variable " + expression.getVariableName() + " to local." - ); + log.warn("Forcing execution mode of variable {} to local.", expression.getVariableName()); expression.setHighestExecutionMode(ExecutionMode.LOCAL); return argument; } @@ -699,12 +698,11 @@ public StaticContext visitValidateTypeExpression(ValidateTypeExpression expressi targetType.getItemType().isObjectItemType() && targetType.getItemType().isCompatibleWithDataFrames(this.configuration) ) { - LogManager.getLogger("ExecutionModeVisitor") - .info( - "Validation against " - + expression.getSequenceType().getItemType().getName() - + " compatible with data frames." - ); + log.info( + "Validation against " + + expression.getSequenceType().getItemType().getName() + + " compatible with data frames." + ); expression.setHighestExecutionMode(DATAFRAMEifConfigurationAllows()); } else { if ( @@ -723,12 +721,11 @@ public StaticContext visitValidateTypeExpression(ValidateTypeExpression expressi targetType.getItemType().isObjectItemType() && targetType.getItemType().isCompatibleWithDataFrames(this.configuration) ) { - LogManager.getLogger("ExecutionModeVisitor") - .info( - "Validation against " - + expression.getSequenceType().getItemType().getName() - + " compatible with data frames." - ); + log.info( + "Validation against " + + expression.getSequenceType().getItemType().getName() + + " compatible with data frames." + ); expression.setHighestExecutionMode(DATAFRAMEifConfigurationAllows()); } else { if ( diff --git a/src/main/java/org/rumbledb/compiler/InferTypeVisitor.java b/src/main/java/org/rumbledb/compiler/InferTypeVisitor.java index 3d4e5bbbc1..79254f3978 100644 --- a/src/main/java/org/rumbledb/compiler/InferTypeVisitor.java +++ b/src/main/java/org/rumbledb/compiler/InferTypeVisitor.java @@ -1,5 +1,7 @@ package org.rumbledb.compiler; +import lombok.extern.log4j.Log4j2; + import java.net.URI; import java.util.ArrayList; import java.util.Arrays; @@ -164,6 +166,7 @@ /** * This visitor infers a static SequenceType for each expression in the query */ +@Log4j2 public class InferTypeVisitor extends AbstractNodeVisitor { private RumbleConfiguration configuration; @@ -393,8 +396,8 @@ public StaticContext visitVariableReference(VariableReferenceExpression expressi variableType = expression.getStaticContext().getVariableSequenceType(expression.getVariableName()); // we also set variableReference type if (variableType == null) { - System.err.println( - "[WARNING] Variable reference type was null so we infer it. Please let us know as we would like to look into it." + log.warn( + "Variable reference type was null so we infer it. Please let us know as we would like to look into it." ); variableType = SequenceType.createSequenceType("item*"); } diff --git a/src/main/java/org/rumbledb/compiler/TranslationVisitor.java b/src/main/java/org/rumbledb/compiler/TranslationVisitor.java index 2e242fd44a..383e14f8fc 100644 --- a/src/main/java/org/rumbledb/compiler/TranslationVisitor.java +++ b/src/main/java/org/rumbledb/compiler/TranslationVisitor.java @@ -20,6 +20,7 @@ package org.rumbledb.compiler; +import lombok.extern.log4j.Log4j2; import org.antlr.v4.runtime.CommonTokenStream; import org.antlr.v4.runtime.ParserRuleContext; import org.antlr.v4.runtime.Token; @@ -210,6 +211,7 @@ * * @author Stefan Irimescu, Can Berker Cikis, Ghislain Fourny, Andrea Rinaldi */ +@Log4j2 public class TranslationVisitor extends JsoniqParserBaseVisitor { private StaticContext moduleContext; @@ -292,7 +294,7 @@ public Node visitMainModule(JsoniqParser.MainModuleContext ctx) { // We override with a context item declaration if not present already. Program program = (Program) this.visitProgram(ctx.program()); if (!prolog.hasContextItemDeclaration() && getExternalVariableType(Name.CONTEXT_ITEM) != null) { - System.err.println("[WARNING] Adding context item declaration."); + log.warn("Adding context item declaration."); prolog.addDeclaration( new VariableDeclaration( Name.CONTEXT_ITEM, @@ -2606,7 +2608,7 @@ public Node visitArrayConstructor(JsoniqParser.ArrayConstructorContext ctx) { createMetadataFromContext(sqCtx) ); } else { - System.err.println("Not concatenating to comma."); + log.debug("Not concatenating to comma."); // In JSONiq 4.0, the square array constructor behaves like in XQuery 4.0. for (JsoniqParser.ExprSingleContext memberCtx : memberCtxs) { memberExpressions.add((Expression) this.visitExprSingle(memberCtx)); diff --git a/src/main/java/org/rumbledb/compiler/VisitorHelpers.java b/src/main/java/org/rumbledb/compiler/VisitorHelpers.java index c08d8e5df3..4a259f1bca 100644 --- a/src/main/java/org/rumbledb/compiler/VisitorHelpers.java +++ b/src/main/java/org/rumbledb/compiler/VisitorHelpers.java @@ -1,5 +1,6 @@ package org.rumbledb.compiler; +import lombok.extern.log4j.Log4j2; import org.antlr.v4.runtime.BailErrorStrategy; import org.antlr.v4.runtime.CharStream; import org.antlr.v4.runtime.CharStreams; @@ -38,6 +39,7 @@ import java.util.ArrayList; import java.util.List; +@Log4j2 public class VisitorHelpers { public static RuntimeIterator generateRuntimeIterator(Node node, RumbleConfiguration conf) { @@ -45,7 +47,7 @@ public static RuntimeIterator generateRuntimeIterator(Node node, RumbleConfigura if (conf.debug().printIteratorTree() || conf.debug().logging()) { StringBuilder sb = new StringBuilder(); result.print(sb, 0); - System.err.println(sb); + log.debug(sb); } return result; } @@ -118,21 +120,35 @@ private static MainModule applyTypeDependentOptimizations(MainModule module) { */ private static void debugPrintTree(Module node, RumbleConfiguration conf) { if (conf.debug().printIteratorTree() || conf.debug().logging()) { - System.err.println("***************"); - System.err.println("Expression tree"); - System.err.println("***************"); - System.err.println("Unset execution modes: " + node.numberOfUnsetExecutionModes()); - System.err.println(node); - System.err.println(); - System.err.println(node.getStaticContext()); + log.debug( + """ + *************** + Expression tree + *************** + Unset execution modes: {} + {} + + {}\ + """, + node.numberOfUnsetExecutionModes(), + node, + node.getStaticContext() + ); } } private static void debugPrintHeader(RumbleConfiguration conf, String header) { if (conf.debug().printIteratorTree() || conf.debug().logging()) { - System.err.println("*".repeat(header.length())); - System.err.println(header); - System.err.println("*".repeat(header.length())); + log.debug( + """ + {} + {} + {}\ + """, + "*".repeat(header.length()), + header, + "*".repeat(header.length()) + ); } } @@ -570,7 +586,7 @@ private static void populateExecutionModes( debugPrintTree(module, conf); } if (module.numberOfUnsetExecutionModes() > 0) { - System.err.println( + log.warn( "[WARNING] Some execution modes could not be set. The query may still work, but we would welcome a bug report." ); } @@ -619,7 +635,7 @@ private static void populateExecutionModes( debugPrintTree(module, conf); } if (module.numberOfUnsetExecutionModes() > 0) { - System.err.println( + log.warn( "[WARNING] Some execution modes could not be set. The query may still work, but we would welcome a bug report." ); } diff --git a/src/main/java/org/rumbledb/compiler/XQueryTranslationVisitor.java b/src/main/java/org/rumbledb/compiler/XQueryTranslationVisitor.java index e8acea0b2a..7398f854d5 100644 --- a/src/main/java/org/rumbledb/compiler/XQueryTranslationVisitor.java +++ b/src/main/java/org/rumbledb/compiler/XQueryTranslationVisitor.java @@ -20,6 +20,7 @@ package org.rumbledb.compiler; +import lombok.extern.log4j.Log4j2; import org.antlr.v4.runtime.CommonTokenStream; import org.antlr.v4.runtime.ParserRuleContext; import org.antlr.v4.runtime.Token; @@ -190,6 +191,7 @@ * * @author Stefan Irimescu, Can Berker Cikis, Ghislain Fourny, Andrea Rinaldi */ +@Log4j2 public class XQueryTranslationVisitor extends XQueryParserBaseVisitor { private StaticContext moduleContext; @@ -276,7 +278,7 @@ public Node visitMainModule(XQueryParser.MainModuleContext ctx) { // We override with a context item declaration if not present already. Program program = (Program) this.visitProgram(ctx.program()); if (!prolog.hasContextItemDeclaration() && getExternalVariableType(Name.CONTEXT_ITEM) != null) { - System.err.println("[WARNING] Adding context item declaration."); + log.warn("Adding context item declaration."); prolog.addDeclaration( new VariableDeclaration( Name.CONTEXT_ITEM, diff --git a/src/main/java/org/rumbledb/config/model/DebugConfig.java b/src/main/java/org/rumbledb/config/model/DebugConfig.java index 82fc7a9b2e..1b798b1bfb 100644 --- a/src/main/java/org/rumbledb/config/model/DebugConfig.java +++ b/src/main/java/org/rumbledb/config/model/DebugConfig.java @@ -67,11 +67,26 @@ public class DebugConfig implements Serializable, KryoSerializable { @Default private boolean logging = false; + /** + * The diagnostic logging level requested by the user. A null value means the CLI startup chooses the default. + */ + @NonFinal + private String logLevel; + + /** + * The Spark logging level requested by the user. + */ + @NonFinal + @Default + private String sparkLogLevel = "off"; + @Override public void write(Kryo kryo, Output output) { output.writeBoolean(this.showErrorInfo); output.writeBoolean(this.printIteratorTree); output.writeBoolean(this.logging); + kryo.writeObjectOrNull(output, this.logLevel, String.class); + kryo.writeObjectOrNull(output, this.sparkLogLevel, String.class); } @Override @@ -79,6 +94,8 @@ public void read(Kryo kryo, Input input) { this.showErrorInfo = input.readBoolean(); this.printIteratorTree = input.readBoolean(); this.logging = input.readBoolean(); + this.logLevel = kryo.readObjectOrNull(input, String.class); + this.sparkLogLevel = kryo.readObjectOrNull(input, String.class); } // Avoid Javadoc error because it cannot resolve the builder class generated by Lombok diff --git a/src/main/java/org/rumbledb/context/DynamicContext.java b/src/main/java/org/rumbledb/context/DynamicContext.java index df4d9c9ca0..572d1a4bf7 100644 --- a/src/main/java/org/rumbledb/context/DynamicContext.java +++ b/src/main/java/org/rumbledb/context/DynamicContext.java @@ -25,7 +25,7 @@ import com.esotericsoftware.kryo.io.Input; import com.esotericsoftware.kryo.io.Output; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.api.java.JavaRDD; import java.time.OffsetDateTime; import org.rumbledb.api.Item; @@ -40,6 +40,8 @@ import java.util.List; import java.util.Map; + +@Log4j2 public class DynamicContext implements Serializable, KryoSerializable { private static final long serialVersionUID = 1L; @@ -284,9 +286,9 @@ public OffsetDateTime getCurrentDateTime() { } public static void printDependencies(Map exprDependency) { - LogManager.getLogger("DynamicContext").debug("System.err Variable dependencies:"); + log.debug("System.err Variable dependencies:"); for (Map.Entry e : exprDependency.entrySet()) { - LogManager.getLogger("DynamicContext").debug(e.getKey() + " : " + e.getValue()); + log.debug(e.getKey() + " : " + e.getValue()); } } diff --git a/src/main/java/org/rumbledb/context/StaticContext.java b/src/main/java/org/rumbledb/context/StaticContext.java index cc1c6a5f4d..bbb9d80a59 100644 --- a/src/main/java/org/rumbledb/context/StaticContext.java +++ b/src/main/java/org/rumbledb/context/StaticContext.java @@ -20,6 +20,8 @@ package org.rumbledb.context; +import lombok.extern.log4j.Log4j2; + import java.io.Serializable; import java.net.URI; import java.util.LinkedHashSet; @@ -47,6 +49,7 @@ import com.esotericsoftware.kryo.io.Input; import com.esotericsoftware.kryo.io.Output; +@Log4j2 public class StaticContext implements Serializable, KryoSerializable { private static final long serialVersionUID = 1L; @@ -326,7 +329,7 @@ public Map getInScopeVariables() { } public void show() { - System.err.println(this); + log.debug(this); } @Override diff --git a/src/main/java/org/rumbledb/expressions/flowr/Clause.java b/src/main/java/org/rumbledb/expressions/flowr/Clause.java index 9c8070b5fa..982629cddb 100644 --- a/src/main/java/org/rumbledb/expressions/flowr/Clause.java +++ b/src/main/java/org/rumbledb/expressions/flowr/Clause.java @@ -20,6 +20,7 @@ package org.rumbledb.expressions.flowr; +import lombok.extern.log4j.Log4j2; import org.rumbledb.compiler.VisitorConfig; import org.rumbledb.config.RumbleConfiguration; @@ -38,6 +39,7 @@ * * Clauses, unlike expressions, return tuple streams. */ +@Log4j2 public abstract class Clause extends Node { /* Clauses are organized in doubly-linked lists */ @@ -125,18 +127,7 @@ public ReturnClause detachInitialLetClauses() { for (Clause c = newFirstClause; c != null; c = c.nextClause) { if (c.getClauseType().equals(FLWOR_CLAUSES.GROUP_BY)) { // No optimization possible if there is a group by. - System.err.println( - "[WARNING] It seems you are using a group by clause in a FLWOR expression that starts with a let clause. This is rather unusual and it might lead to surprises. We recommend always inserting a 'return' after a series of initial let clauses." - ); - System.err.println("For example:"); - System.err.println(); - System.err.println("let $x := 1"); - System.err.println("let $y := $x + 1"); - System.err.println("let $z := $x + $y"); - System.err.println("return"); - System.err.println(" for $t in 1 to $z"); - System.err.println(" group by $m := $t mod 2"); - System.err.println(" return $m + $x"); + logInitialLetGroupByWarning(); return returnClause; } @@ -190,18 +181,7 @@ public ReturnStatementClause detachInitialLetClausesForStatements() { for (Clause c = newFirstClause; c != null; c = c.nextClause) { if (c.getClauseType().equals(FLWOR_CLAUSES.GROUP_BY)) { // No optimization possible if there is a group by. - System.err.println( - "[WARNING] It seems you are using a group by clause in a FLWOR expression that starts with a let clause. This is rather unusual and it might lead to surprises. We recommend always inserting a 'return' after a series of initial let clauses." - ); - System.err.println("For example:"); - System.err.println(); - System.err.println("let $x := 1"); - System.err.println("let $y := $x + 1"); - System.err.println("let $z := $x + $y"); - System.err.println("return"); - System.err.println(" for $t in 1 to $z"); - System.err.println(" group by $m := $t mod 2"); - System.err.println(" return $m + $x"); + logInitialLetGroupByWarning(); return returnClause; } @@ -222,6 +202,23 @@ public ReturnStatementClause detachInitialLetClausesForStatements() { return returnClause; } + private static void logInitialLetGroupByWarning() { + log.warn( + """ + It seems you are using a group by clause in a FLWOR expression that starts with a let clause. This is rather unusual and it might lead to surprises. We recommend always inserting a 'return' after a series of initial let clauses. + For example: + + let $x := 1 + let $y := $x + 1 + let $z := $x + $y + return + for $t in 1 to $z + group by $m := $t mod 2 + return $m + $x\ + """ + ); + } + public void print(StringBuilder buffer, int indent) { for (int i = 0; i < indent; ++i) { diff --git a/src/main/java/org/rumbledb/items/structured/JSoundDataFrame.java b/src/main/java/org/rumbledb/items/structured/JSoundDataFrame.java index 1862432e13..43ebe41211 100644 --- a/src/main/java/org/rumbledb/items/structured/JSoundDataFrame.java +++ b/src/main/java/org/rumbledb/items/structured/JSoundDataFrame.java @@ -1,5 +1,7 @@ package org.rumbledb.items.structured; +import lombok.extern.log4j.Log4j2; + import java.io.Serializable; import java.util.ArrayList; import java.util.Arrays; @@ -23,6 +25,7 @@ import org.rumbledb.types.ItemType; import org.rumbledb.types.ItemTypeFactory; +@Log4j2 public class JSoundDataFrame implements Serializable { private static final long serialVersionUID = 1L; @@ -128,7 +131,7 @@ public Dataset getDataFrame() { } public void show() { - System.out.println("Item type: " + this.itemType); + log.debug("Item type: {}", this.itemType); this.dataFrame.show(); } diff --git a/src/main/java/org/rumbledb/optimizations/Profiler.java b/src/main/java/org/rumbledb/optimizations/Profiler.java index a650f7c298..5c8cbdce66 100644 --- a/src/main/java/org/rumbledb/optimizations/Profiler.java +++ b/src/main/java/org/rumbledb/optimizations/Profiler.java @@ -1,11 +1,14 @@ package org.rumbledb.optimizations; +import lombok.extern.log4j.Log4j2; + import java.util.Arrays; import java.util.Collections; import java.util.HashMap; import java.util.Map; import java.util.stream.Collectors; +@Log4j2 public class Profiler { public static int counter = 0; @@ -37,12 +40,17 @@ public static int get() { total += stacks.get(key); if (stacks.get(key) != max) continue; - System.err.println("Occurrences: " + stacks.get(key)); - System.err.println(key); - System.err.println(); + log.debug( + """ + Occurrences: {} + {}\ + """, + stacks.get(key), + key + ); } - System.err.println("Size: " + stacks.size()); - System.err.println("Total: " + total); + log.debug("Size: {}", stacks.size()); + log.debug("Total: {}", total); return counter; } diff --git a/src/main/java/org/rumbledb/runtime/AtMostOneItemLocalRuntimeIterator.java b/src/main/java/org/rumbledb/runtime/AtMostOneItemLocalRuntimeIterator.java index ae8faec32b..f474d20743 100644 --- a/src/main/java/org/rumbledb/runtime/AtMostOneItemLocalRuntimeIterator.java +++ b/src/main/java/org/rumbledb/runtime/AtMostOneItemLocalRuntimeIterator.java @@ -20,6 +20,8 @@ package org.rumbledb.runtime; +import lombok.extern.log4j.Log4j2; + import org.apache.spark.api.java.JavaRDD; import org.rumbledb.api.Item; import org.rumbledb.context.DynamicContext; @@ -40,6 +42,7 @@ import java.util.ArrayList; import java.util.List; +@Log4j2 public abstract class AtMostOneItemLocalRuntimeIterator extends RuntimeIterator { private static final long serialVersionUID = 1L; @@ -180,10 +183,12 @@ public boolean getEffectiveBooleanValueOrCheckPosition(DynamicContext dynamicCon } } else { if (item.isObject() || item.isArray()) { - System.err.println( - "Note: effective boolean value of " - + (item.isObject() ? "Object " : "Array ") - + "accessed which throws error in JSONiq 3.1 or 4.0 in alignment with Xquery 3.1 or 4.0 spec.\n If you want to revert to the old functionality use the --default-language jsoniq10 command line option" + log.warn( + """ + Note: effective boolean value of {} accessed which throws error in JSONiq 3.1 or 4.0 in alignment with Xquery 3.1 or 4.0 spec. + If you want to revert to the old functionality use the --default-language jsoniq10 command line option\ + """, + item.isObject() ? "Object" : "Array" ); } } diff --git a/src/main/java/org/rumbledb/runtime/RuntimeIterator.java b/src/main/java/org/rumbledb/runtime/RuntimeIterator.java index 323d746caf..fe7071b3b9 100644 --- a/src/main/java/org/rumbledb/runtime/RuntimeIterator.java +++ b/src/main/java/org/rumbledb/runtime/RuntimeIterator.java @@ -20,6 +20,8 @@ package org.rumbledb.runtime; +import lombok.extern.log4j.Log4j2; + import java.io.ByteArrayInputStream; import java.io.ByteArrayOutputStream; import java.io.IOException; @@ -66,6 +68,7 @@ import com.esotericsoftware.kryo.io.Input; import com.esotericsoftware.kryo.io.Output; +@Log4j2 public abstract class RuntimeIterator implements RuntimeIteratorInterface, KryoSerializable { protected static final String FLOW_EXCEPTION_MESSAGE = "Invalid next() call; "; @@ -168,10 +171,12 @@ public boolean getEffectiveBooleanValueOrCheckPosition(DynamicContext dynamicCon } } else { if (item.isObject() || item.isArray()) { - System.err.println( - "Note: effective boolean value of " - + (item.isObject() ? "Object " : "Array ") - + "accessed which throws error in JSONiq 3.1 or 4.0 in alignment with Xquery 3.1 or 4.0 spec.\n If you want to revert to the old functionality use the --default-language jsoniq10 command line option" + log.warn( + """ + Note: effective boolean value of {} accessed which throws error in JSONiq 3.1 or 4.0 in alignment with Xquery 3.1 or 4.0 spec. + If you want to revert to the old functionality use the --default-language jsoniq10 command line option\ + """, + item.isObject() ? "Object" : "Array" ); } } @@ -393,7 +398,13 @@ public final JSoundDataFrame getOrCreateDataFrame(DynamicContext context) { TypeInferrenceUtils.TypeMergeMode.LAX ); if (this.getConfiguration().analysis().printInferredTypes()) { - System.err.println("Inferred DataFrame type:\n" + this.getStaticType().getItemType()); + log.debug( + """ + Inferred DataFrame type: + {}\ + """, + this.getStaticType().getItemType() + ); } return ValidateTypeIterator.convertLocalItemsToDataFrame( items, @@ -531,7 +542,7 @@ public Map getVariableDependencies() { public void printToStandardError() { StringBuilder sb = new StringBuilder(); this.print(sb, 0); - System.err.println(sb); + log.debug(sb); } public void print(StringBuilder buffer, int indent) { diff --git a/src/main/java/org/rumbledb/runtime/flwor/FlworDataFrame.java b/src/main/java/org/rumbledb/runtime/flwor/FlworDataFrame.java index f92c70df81..dd4af8dd29 100644 --- a/src/main/java/org/rumbledb/runtime/flwor/FlworDataFrame.java +++ b/src/main/java/org/rumbledb/runtime/flwor/FlworDataFrame.java @@ -1,5 +1,7 @@ package org.rumbledb.runtime.flwor; +import lombok.extern.log4j.Log4j2; + import java.io.Serializable; import java.util.ArrayList; import java.util.HashMap; @@ -15,6 +17,7 @@ import org.rumbledb.exceptions.OurBadException; import org.rumbledb.types.SequenceType; +@Log4j2 public class FlworDataFrame implements Serializable { private static final long serialVersionUID = 1L; @@ -117,16 +120,18 @@ public String createTempView() { } public void show() { - System.err.println("FLWOR DataFrame"); - System.err.println("Columns"); + StringBuilder sb = new StringBuilder(); + sb.append("FLWOR DataFrame\n"); + sb.append("Columns\n"); for (FlworDataFrameColumn c : this.columns) { - System.err.println(c.toString()); + sb.append(c).append('\n'); } - System.err.println("Column types"); + sb.append("Column types\n"); for (Name n : this.columnTypes.keySet()) { - System.err.println(n + " " + this.columnTypes.get(n)); + sb.append(n).append(' ').append(this.columnTypes.get(n)).append('\n'); } - System.err.println("Data Frame"); + sb.append("Data Frame"); + log.debug(sb); this.dataFrame.show(); } diff --git a/src/main/java/org/rumbledb/runtime/flwor/clauses/ForClauseIterator.java b/src/main/java/org/rumbledb/runtime/flwor/clauses/ForClauseIterator.java index d67b46b04c..9fbda14364 100644 --- a/src/main/java/org/rumbledb/runtime/flwor/clauses/ForClauseIterator.java +++ b/src/main/java/org/rumbledb/runtime/flwor/clauses/ForClauseIterator.java @@ -20,7 +20,7 @@ package org.rumbledb.runtime.flwor.clauses; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.api.java.JavaRDD; import org.apache.spark.api.java.JavaSparkContext; import org.apache.spark.sql.Dataset; @@ -73,6 +73,8 @@ import java.util.TreeMap; + +@Log4j2 public class ForClauseIterator extends RuntimeTupleIterator { @@ -1022,10 +1024,9 @@ public static FlworDataFrame tryNativeQuery( if (nativeQuery == NativeClauseContext.NoNativeQuery) { return null; } - LogManager.getLogger("ForClauseSparkIterator") - .info( - "Rumble was able to optimize a for clause to a native SQL query." - ); + log.info( + "Rumble was able to optimize a for clause to a native SQL query." + ); String selectSQL = FlworDataFrameUtils.getSQLColumnProjection(allColumns, true); String viewName = FlworDataFrameUtils.createTempView(dataFrame); diff --git a/src/main/java/org/rumbledb/runtime/flwor/clauses/GroupByClauseIterator.java b/src/main/java/org/rumbledb/runtime/flwor/clauses/GroupByClauseIterator.java index 62ff88bc6e..24adc08978 100644 --- a/src/main/java/org/rumbledb/runtime/flwor/clauses/GroupByClauseIterator.java +++ b/src/main/java/org/rumbledb/runtime/flwor/clauses/GroupByClauseIterator.java @@ -20,7 +20,7 @@ package org.rumbledb.runtime.flwor.clauses; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.types.DataType; @@ -67,6 +67,8 @@ import java.util.TreeMap; import java.util.stream.Collectors; + +@Log4j2 public class GroupByClauseIterator extends RuntimeTupleIterator { private static final long serialVersionUID = 1L; @@ -667,8 +669,7 @@ private Dataset tryNativeQuery( selectString.append(") as "); selectString.append(dfColumnSequence); } - LogManager.getLogger("GroupByClauseSparkIterator") - .info("Rumble was able to optimize a group by clause to a native SQL query."); + log.info("Rumble was able to optimize a group by clause to a native SQL query."); return dataFrame.sparkSession() .sql( String.format( diff --git a/src/main/java/org/rumbledb/runtime/flwor/clauses/JoinClauseIterator.java b/src/main/java/org/rumbledb/runtime/flwor/clauses/JoinClauseIterator.java index fe519f64f9..df333b1157 100644 --- a/src/main/java/org/rumbledb/runtime/flwor/clauses/JoinClauseIterator.java +++ b/src/main/java/org/rumbledb/runtime/flwor/clauses/JoinClauseIterator.java @@ -28,7 +28,7 @@ import java.util.Set; import java.util.Stack; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.types.DataTypes; @@ -57,6 +57,8 @@ import org.rumbledb.types.SequenceType; + +@Log4j2 public class JoinClauseIterator extends RuntimeTupleIterator { private static final long serialVersionUID = 1L; @@ -187,10 +189,9 @@ public static FlworDataFrame joinInputTupleWithSequenceOnPredicate( // } if (optimizableJoin) { - LogManager.getLogger("JoinClauseSparkIterator") - .info( - "Rumble detected that it can optimize your query and make it faster with an equi-join." - ); + log.info( + "Rumble detected that it can optimize your query and make it faster with an equi-join." + ); } @@ -473,10 +474,9 @@ private static FlworDataFrame tryNativeQueryStatically( if (nativeQuery == NativeClauseContext.NoNativeQuery) { return null; } - LogManager.getLogger("JoinClauseSparkIterator") - .info( - "Rumble was able to optimize a join to a native SQL query." - ); + log.info( + "Rumble was able to optimize a join to a native SQL query." + ); String left = FlworDataFrameUtils.createTempView(leftInputTuple); String right = FlworDataFrameUtils.createTempView(rightInputTuple); List columnsToSelect = FlworDataFrameUtils.getColumns( diff --git a/src/main/java/org/rumbledb/runtime/flwor/clauses/LetClauseIterator.java b/src/main/java/org/rumbledb/runtime/flwor/clauses/LetClauseIterator.java index b30468e9e6..a1cc210de5 100644 --- a/src/main/java/org/rumbledb/runtime/flwor/clauses/LetClauseIterator.java +++ b/src/main/java/org/rumbledb/runtime/flwor/clauses/LetClauseIterator.java @@ -20,7 +20,7 @@ package org.rumbledb.runtime.flwor.clauses; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.api.java.JavaRDD; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; @@ -66,6 +66,8 @@ import java.math.BigDecimal; import java.util.*; + +@Log4j2 public class LetClauseIterator extends RuntimeTupleIterator { @@ -871,19 +873,18 @@ public static Dataset tryNativeQuery( return null; } String selectSQL = FlworDataFrameUtils.getSQLColumnProjection(allColumns, true); - LogManager.getLogger("LetClauseSparkIterator") - .info( - "Rumble was able to optimize a let clause to a native SQL query: " - + String.format( - "select %s %s as `%s` from (%s)", - selectSQL, - nativeQuery.getResultingQuery(), - SequenceType.Arity.OneOrMore.isSubtypeOf(nativeQuery.getResultingType().getArity()) - ? newVariableName + ".sequence" - : newVariableName, - nativeQuery.getView() - ) - ); + log.info( + "Rumble was able to optimize a let clause to a native SQL query: " + + String.format( + "select %s %s as `%s` from (%s)", + selectSQL, + nativeQuery.getResultingQuery(), + SequenceType.Arity.OneOrMore.isSubtypeOf(nativeQuery.getResultingType().getArity()) + ? newVariableName + ".sequence" + : newVariableName, + nativeQuery.getView() + ) + ); return dataFrame.sparkSession() .sql( String.format( diff --git a/src/main/java/org/rumbledb/runtime/flwor/clauses/OrderByClauseIterator.java b/src/main/java/org/rumbledb/runtime/flwor/clauses/OrderByClauseIterator.java index 7ba2f70836..8b93fe8c75 100644 --- a/src/main/java/org/rumbledb/runtime/flwor/clauses/OrderByClauseIterator.java +++ b/src/main/java/org/rumbledb/runtime/flwor/clauses/OrderByClauseIterator.java @@ -20,7 +20,7 @@ package org.rumbledb.runtime.flwor.clauses; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.types.DataType; @@ -57,6 +57,8 @@ import java.util.*; + +@Log4j2 public class OrderByClauseIterator extends RuntimeTupleIterator { public static final String StringFlagForEmptySequence = "empty-sequence"; @@ -531,9 +533,7 @@ public static FlworDataFrame tryNativeQuery( NativeClauseContext queryContext = createOrderExpression(expressionsWithIterator, orderContext); if (queryContext == NativeClauseContext.NoNativeQuery) return null; - - LogManager.getLogger("OrderByClauseSparkIterator") - .info("Rumble was able to optimize an order-by clause to a native SQL query."); + log.info("Rumble was able to optimize an order-by clause to a native SQL query."); String selectSQL = FlworDataFrameUtils.getSQLColumnProjection(allColumns, false); dataFrame.createOrReplaceTempView("input"); return new FlworDataFrame( diff --git a/src/main/java/org/rumbledb/runtime/flwor/clauses/ReturnClauseIterator.java b/src/main/java/org/rumbledb/runtime/flwor/clauses/ReturnClauseIterator.java index c672342a5d..ff7a0c9460 100644 --- a/src/main/java/org/rumbledb/runtime/flwor/clauses/ReturnClauseIterator.java +++ b/src/main/java/org/rumbledb/runtime/flwor/clauses/ReturnClauseIterator.java @@ -20,7 +20,7 @@ package org.rumbledb.runtime.flwor.clauses; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.api.java.JavaRDD; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; @@ -63,6 +63,8 @@ import java.util.TreeMap; import java.util.stream.Collectors; + +@Log4j2 public class ReturnClauseIterator extends HybridRuntimeIterator { private static final long serialVersionUID = 1L; @@ -396,11 +398,10 @@ public static Dataset tryNativeQuery( SparkSessionManager.nonObjectJSONiqItemColumnName ); } - LogManager.getLogger("ReturnClauseSparkIterator") - .info( - "Rumble was able to optimize a return clause to a native SQL query: " - + queryString - ); + log.info( + "Rumble was able to optimize a return clause to a native SQL query: " + + queryString + ); return dataFrame.sparkSession().sql(queryString); } diff --git a/src/main/java/org/rumbledb/runtime/flwor/clauses/WhereClauseIterator.java b/src/main/java/org/rumbledb/runtime/flwor/clauses/WhereClauseIterator.java index 473c76bfa2..b00785f9a8 100644 --- a/src/main/java/org/rumbledb/runtime/flwor/clauses/WhereClauseIterator.java +++ b/src/main/java/org/rumbledb/runtime/flwor/clauses/WhereClauseIterator.java @@ -20,7 +20,7 @@ package org.rumbledb.runtime.flwor.clauses; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.sql.types.DataTypes; import org.apache.spark.sql.types.StructType; import org.rumbledb.api.Item; @@ -50,6 +50,8 @@ import java.util.*; + +@Log4j2 public class WhereClauseIterator extends RuntimeTupleIterator { @@ -248,10 +250,9 @@ private FlworDataFrame getDataFrameIfLimit(DynamicContext context) { if (!item.isInteger()) { return null; } - LogManager.getLogger("WhereClauseSparkIterator") - .info( - "Rumble detected a LIMIT in a count and where clause." - ); + log.info( + "Rumble detected a LIMIT in a count and where clause." + ); FlworDataFrame df = this.child.getChildIterator().getDataFrame(context); String input = df.createTempView(); return df.sql(String.format("SELECT * FROM %s LIMIT %s", input, item.getStringValue())); @@ -318,9 +319,7 @@ private FlworDataFrame getDataFrameIfJoinPossible(DynamicContext context) { if (limit == -1) { return null; } - - LogManager.getLogger("WhereClauseSparkIterator") - .info("Rumble detected a join predicate in the where clause (limit=" + limit + " of " + height + ")."); + log.info("Rumble detected a join predicate in the where clause (limit=" + limit + " of " + height + ")."); try { FlworDataFrame leftTuples = getSubtreeBeyondLimit(limit).getDataFrame(context); @@ -352,10 +351,9 @@ private FlworDataFrame getDataFrameIfJoinPossible(DynamicContext context) { ); return result; } catch (Exception e) { - LogManager.getLogger("WhereClauseSparkIterator") - .warn( - "Join failed. Falling back to regular execution (nevertheless, please let us know!)." - ); + log.warn( + "Join failed. Falling back to regular execution (nevertheless, please let us know!)." + ); this.setEvaluationDepthLimit(-1); this.child.setInputAndOutputTupleVariableDependencies(this.inputTupleProjection); @@ -446,16 +444,15 @@ public static FlworDataFrame tryNativeQuery( iterator.getMetadata() ); } - LogManager.getLogger("WhereClauseSparkIterator") - .info( - "Rumble was able to optimize a where clause to a native SQL query: " - + String.format( - "select %s from (%s) where %s", - FlworDataFrameUtils.getSQLColumnProjection(allColumns, false), - nativeQuery.getView(), - nativeQuery.getResultingQuery() - ) - ); + log.info( + "Rumble was able to optimize a where clause to a native SQL query: " + + String.format( + "select %s from (%s) where %s", + FlworDataFrameUtils.getSQLColumnProjection(allColumns, false), + nativeQuery.getView(), + nativeQuery.getResultingQuery() + ) + ); return new FlworDataFrame( dataFrame.getDataFrame() .sparkSession() diff --git a/src/main/java/org/rumbledb/runtime/flwor/expression/SimpleMapExpressionIterator.java b/src/main/java/org/rumbledb/runtime/flwor/expression/SimpleMapExpressionIterator.java index eb4bbc9d10..d1cd5a462f 100644 --- a/src/main/java/org/rumbledb/runtime/flwor/expression/SimpleMapExpressionIterator.java +++ b/src/main/java/org/rumbledb/runtime/flwor/expression/SimpleMapExpressionIterator.java @@ -20,7 +20,7 @@ package org.rumbledb.runtime.flwor.expression; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.api.java.JavaPairRDD; import org.apache.spark.api.java.JavaRDD; import org.apache.spark.api.java.function.FlatMapFunction; @@ -53,6 +53,8 @@ import java.util.Queue; import java.util.TreeMap; + +@Log4j2 public class SimpleMapExpressionIterator extends HybridRuntimeIterator { private static final long serialVersionUID = 1L; @@ -214,8 +216,7 @@ public JSoundDataFrame getDataFrame(DynamicContext context) { .createDataFrame(rowRDD, schema); return new JSoundDataFrame(result, getStaticType().getItemType()); } - LogManager.getLogger("SimpleMapExpressionIterator") - .info("Rumble was able to optimize a simple map expression to a native SQL query."); + log.info("Rumble was able to optimize a simple map expression to a native SQL query."); String input = FlworDataFrameUtils.createTempView(df.getDataFrame()); Dataset result = df.getDataFrame() .sparkSession() diff --git a/src/main/java/org/rumbledb/runtime/functions/sequences/general/ReverseFunctionIterator.java b/src/main/java/org/rumbledb/runtime/functions/sequences/general/ReverseFunctionIterator.java index 8262dd3b97..74ce06c1da 100644 --- a/src/main/java/org/rumbledb/runtime/functions/sequences/general/ReverseFunctionIterator.java +++ b/src/main/java/org/rumbledb/runtime/functions/sequences/general/ReverseFunctionIterator.java @@ -20,7 +20,7 @@ package org.rumbledb.runtime.functions.sequences.general; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.api.java.JavaPairRDD; import org.apache.spark.api.java.JavaRDD; import org.rumbledb.api.Item; @@ -38,6 +38,8 @@ import java.util.ArrayList; import java.util.List; + +@Log4j2 public class ReverseFunctionIterator extends HybridRuntimeIterator { @@ -71,17 +73,16 @@ public JSoundDataFrame getDataFrame(DynamicContext context) { JSoundDataFrame childDataFrame = this.children.get(0).getDataFrame(context); String viewName = FlworDataFrameUtils.createTempView(childDataFrame.getDataFrame()); String selectSQL = childDataFrame.getSQLColumnProjection(false); - LogManager.getLogger("ReverseFunctioniterator") - .info( - String.format( - "SELECT %s FROM (SELECT %s, monotonically_increasing_id() as `%s` FROM %s ORDER BY `%s` DESC)", - selectSQL, - selectSQL, - "foo", - viewName, - "foo" - ) - ); + log.info( + String.format( + "SELECT %s FROM (SELECT %s, monotonically_increasing_id() as `%s` FROM %s ORDER BY `%s` DESC)", + selectSQL, + selectSQL, + "foo", + viewName, + "foo" + ) + ); String tempName = SparkSessionManager.temporaryColumnName; JSoundDataFrame result = childDataFrame.evaluateSQL( String.format( diff --git a/src/main/java/org/rumbledb/runtime/functions/typing/FunctionNameFunctionIterator.java b/src/main/java/org/rumbledb/runtime/functions/typing/FunctionNameFunctionIterator.java index cea5b5765b..24d166a5ca 100644 --- a/src/main/java/org/rumbledb/runtime/functions/typing/FunctionNameFunctionIterator.java +++ b/src/main/java/org/rumbledb/runtime/functions/typing/FunctionNameFunctionIterator.java @@ -1,5 +1,7 @@ package org.rumbledb.runtime.functions.typing; +import lombok.extern.log4j.Log4j2; + import org.rumbledb.api.Item; import org.rumbledb.context.DynamicContext; import org.rumbledb.context.Name; @@ -14,6 +16,7 @@ import java.util.List; +@Log4j2 public class FunctionNameFunctionIterator extends AtMostOneItemLocalRuntimeIterator { private static final long serialVersionUID = 1L; @@ -37,7 +40,7 @@ public Item materializeFirstItemOrNull(DynamicContext context) { getMetadata() ); } - System.err.println("Item is of type function"); + log.debug("Item is of type function"); Item functionItem = functionIterator.materializeFirstItemOrNull(context); if (!(functionItem instanceof FunctionItem function)) { throw new OurBadException("Expected argument to be of type function and not be null"); diff --git a/src/main/java/org/rumbledb/runtime/navigation/ArrayLookupIterator.java b/src/main/java/org/rumbledb/runtime/navigation/ArrayLookupIterator.java index 110fadbfd6..d52851b6ea 100644 --- a/src/main/java/org/rumbledb/runtime/navigation/ArrayLookupIterator.java +++ b/src/main/java/org/rumbledb/runtime/navigation/ArrayLookupIterator.java @@ -20,7 +20,7 @@ package org.rumbledb.runtime.navigation; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.api.java.JavaRDD; import org.apache.spark.api.java.function.FlatMapFunction; import org.apache.spark.sql.Dataset; @@ -47,6 +47,8 @@ import java.util.Arrays; import java.util.Map; + +@Log4j2 public class ArrayLookupIterator extends HybridRuntimeIterator { @@ -221,10 +223,9 @@ public NativeClauseContext generateNativeQuery(NativeClauseContext nativeClauseC getMetadata() ); } - LogManager.getLogger("ArrayLookupIterator") - .warn( - "Array lookup on a DataFrame that does not an array type. Empty sequence returned." - ); + log.warn( + "Array lookup on a DataFrame that does not an array type. Empty sequence returned." + ); return NativeClauseContext.NoNativeQuery; } @@ -239,10 +240,9 @@ public NativeClauseContext generateNativeQuery(NativeClauseContext nativeClauseC getMetadata() ); } - LogManager.getLogger("ArrayLookupIterator") - .warn( - "Array lookup on a DataFrame that does not an array type. Empty sequence returned." - ); + log.warn( + "Array lookup on a DataFrame that does not an array type. Empty sequence returned." + ); return NativeClauseContext.NoNativeQuery; } newContext.setResultingType( @@ -402,10 +402,9 @@ public JSoundDataFrame getDataFrame(DynamicContext context) { getMetadata() ); } - LogManager.getLogger("ArrayLookupIterator") - .warn( - "Array lookup on a DataFrame that does not an array type. Empty sequence returned." - ); + log.warn( + "Array lookup on a DataFrame that does not an array type. Empty sequence returned." + ); return JSoundDataFrame.emptyDataFrame(); } } diff --git a/src/main/java/org/rumbledb/runtime/navigation/ArrayUnboxingIterator.java b/src/main/java/org/rumbledb/runtime/navigation/ArrayUnboxingIterator.java index 3b482fe0f2..8c8a574b9d 100644 --- a/src/main/java/org/rumbledb/runtime/navigation/ArrayUnboxingIterator.java +++ b/src/main/java/org/rumbledb/runtime/navigation/ArrayUnboxingIterator.java @@ -20,7 +20,7 @@ package org.rumbledb.runtime.navigation; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.api.java.JavaRDD; import org.apache.spark.api.java.function.FlatMapFunction; import org.apache.spark.sql.Dataset; @@ -48,6 +48,8 @@ import java.util.List; import java.util.Queue; + +@Log4j2 public class ArrayUnboxingIterator extends HybridRuntimeIterator { private static final long serialVersionUID = 1L; @@ -157,10 +159,9 @@ public NativeClauseContext generateNativeQuery(NativeClauseContext nativeClauseC getMetadata() ); } - LogManager.getLogger("ArrayUnboxingIterator") - .warn( - "Array unboxing on a DataFrame that does not an array type. Empty sequence returned." - ); + log.warn( + "Array unboxing on a DataFrame that does not an array type. Empty sequence returned." + ); return NativeClauseContext.NoNativeQuery; } newContext.setResultingType( @@ -333,10 +334,9 @@ public JSoundDataFrame getDataFrame(DynamicContext context) { getMetadata() ); } - LogManager.getLogger("ArrayUnboxingIterator") - .warn( - "Array unboxing on a DataFrame that does not an array type. Empty sequence returned." - ); + log.warn( + "Array unboxing on a DataFrame that does not an array type. Empty sequence returned." + ); return JSoundDataFrame.emptyDataFrame(); } } diff --git a/src/main/java/org/rumbledb/runtime/navigation/ObjectLookupIterator.java b/src/main/java/org/rumbledb/runtime/navigation/ObjectLookupIterator.java index af0d8640df..822587048c 100644 --- a/src/main/java/org/rumbledb/runtime/navigation/ObjectLookupIterator.java +++ b/src/main/java/org/rumbledb/runtime/navigation/ObjectLookupIterator.java @@ -25,7 +25,7 @@ import java.util.Arrays; import java.util.Map; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.api.java.JavaRDD; import org.apache.spark.api.java.function.FlatMapFunction; import org.apache.spark.sql.types.ArrayType; @@ -61,6 +61,8 @@ import org.rumbledb.types.SequenceType; import org.rumbledb.types.TypeMappings; + +@Log4j2 public class ObjectLookupIterator extends HybridRuntimeIterator { private static final long serialVersionUID = 1L; @@ -308,10 +310,9 @@ public NativeClauseContext generateNativeQuery(NativeClauseContext nativeClauseC getMetadata() ); } - LogManager.getLogger("ObjectLookupIterator") - .warn( - "Object lookup on a DataFrame that does not have this column. Empty sequence returned." - ); + log.warn( + "Object lookup on a DataFrame that does not have this column. Empty sequence returned." + ); } return NativeClauseContext.NoNativeQuery; } @@ -365,10 +366,9 @@ public NativeClauseContext generateNativeQuery(NativeClauseContext nativeClauseC newContext.setSchema(field.dataType()); } else { if (this.children.get(1) instanceof StringRuntimeIterator) { - LogManager.getLogger("ObjectLookupIterator") - .warn( - "Object lookup on a DataFrame that does not have this column. Empty sequence returned." - ); + log.warn( + "Object lookup on a DataFrame that does not have this column. Empty sequence returned." + ); if (getConfiguration().analysis().enableStaticTyping()) { throw new UnexpectedStaticTypeException( "There is no field with the name " @@ -462,10 +462,9 @@ public JSoundDataFrame getDataFrame(DynamicContext context) { return result; } } - LogManager.getLogger("ObjectLookupIterator") - .warn( - "Object lookup on a DataFrame that does not have this column. Empty sequence returned." - ); + log.warn( + "Object lookup on a DataFrame that does not have this column. Empty sequence returned." + ); JSoundDataFrame result = JSoundDataFrame.emptyDataFrame(); return result; } diff --git a/src/main/java/org/rumbledb/runtime/navigation/PredicateIterator.java b/src/main/java/org/rumbledb/runtime/navigation/PredicateIterator.java index b958269e86..ff20abf077 100644 --- a/src/main/java/org/rumbledb/runtime/navigation/PredicateIterator.java +++ b/src/main/java/org/rumbledb/runtime/navigation/PredicateIterator.java @@ -20,7 +20,7 @@ package org.rumbledb.runtime.navigation; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.api.java.JavaPairRDD; import org.apache.spark.api.java.JavaRDD; import org.apache.spark.api.java.function.Function; @@ -57,6 +57,8 @@ import java.math.BigDecimal; import java.util.*; + +@Log4j2 public class PredicateIterator extends HybridRuntimeIterator { private static final long serialVersionUID = 1L; @@ -346,10 +348,9 @@ public JSoundDataFrame getDataFrame(DynamicContext context) { ); } } - LogManager.getLogger("PredicateIterator") - .info( - "Rumble was able to optimize a predicate to a native SQL query." - ); + log.info( + "Rumble was able to optimize a predicate to a native SQL query." + ); String left = FlworDataFrameUtils.createTempView(childDataFrame.getDataFrame()); return childDataFrame.evaluateSQL( String.format( diff --git a/src/main/java/org/rumbledb/runtime/typing/ValidateTypeIterator.java b/src/main/java/org/rumbledb/runtime/typing/ValidateTypeIterator.java index 15175af189..97492b1221 100644 --- a/src/main/java/org/rumbledb/runtime/typing/ValidateTypeIterator.java +++ b/src/main/java/org/rumbledb/runtime/typing/ValidateTypeIterator.java @@ -1,5 +1,7 @@ package org.rumbledb.runtime.typing; +import lombok.extern.log4j.Log4j2; + import java.sql.Date; import java.sql.Timestamp; import java.util.ArrayList; @@ -40,6 +42,7 @@ import org.rumbledb.types.ItemTypeFactory; import org.rumbledb.types.TypeMappings; +@Log4j2 public class ValidateTypeIterator extends HybridRuntimeIterator { private static final long serialVersionUID = 1L; @@ -269,8 +272,13 @@ public static JSoundDataFrame convertLocalItemsToDataFrame( } StructType schema = convertToDataFrameSchema(itemType, staticContext); if (staticContext.getConfiguration().analysis().printInferredTypes()) { - System.err.println("Inferred DataFrame type:\n"); - schema.printTreeString(); + log.debug( + """ + Inferred DataFrame type: + {}\ + """, + schema.treeString() + ); } List rows = new ArrayList<>(); for (Item item : items) { diff --git a/src/main/java/org/rumbledb/spark/SparkSessionManager.java b/src/main/java/org/rumbledb/spark/SparkSessionManager.java index 06aca3d2ba..3db9f3d4e3 100644 --- a/src/main/java/org/rumbledb/spark/SparkSessionManager.java +++ b/src/main/java/org/rumbledb/spark/SparkSessionManager.java @@ -21,7 +21,7 @@ package org.rumbledb.spark; import org.apache.logging.log4j.Level; -import org.apache.logging.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.parquet.format.IntType; import org.apache.spark.SparkConf; import org.apache.spark.api.java.JavaSparkContext; @@ -65,12 +65,14 @@ import org.rumbledb.types.ItemType; import org.rumbledb.types.SequenceType; + +@Log4j2 public class SparkSessionManager { private static final String APP_NAME = "Rumble application"; private static final String DEFAULT_APP_NAME = ""; private static SparkSessionManager instance; - private static final Level LOG_LEVEL = Level.FATAL; + private static Level LOG_LEVEL = Level.OFF; private SparkConf configuration; private SparkSession session; private JavaSparkContext javaSparkContext; @@ -134,6 +136,16 @@ public static SparkSessionManager getInstance(SparkSession session) { return instance; } + public static void setLogLevel(Level level) { + LOG_LEVEL = level; + if (instance != null && instance.session != null) { + instance.session.sparkContext().setLogLevel(level.name()); + } + if (instance != null && instance.configuration != null) { + instance.applySparkLogLevelToConfiguration(); + } + } + public SparkSession getOrCreateSession() { if (this.configuration == null) { setDefaultConfiguration(); @@ -148,11 +160,10 @@ private void setDefaultConfiguration() { try { this.configuration = new SparkConf(); if (this.configuration.get("spark.app.name", DEFAULT_APP_NAME).equals(DEFAULT_APP_NAME)) { - LogManager.getLogger("SparkSessionManager") - .warn( - "No app name specified (you can do so with --conf spark.app.name=your_name). Setting to " - + APP_NAME - ); + log.warn( + "No app name specified (you can do so with --conf spark.app.name=your_name). Setting to " + + APP_NAME + ); this.configuration.setAppName(APP_NAME); } this.configuration.set("spark.mongodb.read.connection.uri", "mongodb://127.0.0.1/test.myCollection"); @@ -165,7 +176,7 @@ private void setDefaultConfiguration() { if (!this.configuration.contains("spark.master")) { this.configuration.set("spark.master", "local[*]"); } - this.configuration.set("spark.log.level", LOG_LEVEL.name()); + applySparkLogLevelToConfiguration(); } catch (NoClassDefFoundError e) { throw new RuntimeException( "It seems your query needs Spark, but it is not available. You need to use spark-submit in an environment in which Spark is configured." @@ -185,12 +196,17 @@ public void resetSession() { private void initializeSession() { if (this.session == null) { initializeKryoSerialization(); + this.session = SparkSession.builder().config(this.configuration).enableHiveSupport().getOrCreate(); } else { throw new OurBadException("Session already exists: new session initialization prevented."); } } + private void applySparkLogLevelToConfiguration() { + this.configuration.set("spark.log.level", LOG_LEVEL.name()); + } + private void initializeKryoSerialization() { if (!this.configuration.contains("spark.serializer")) { this.configuration.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer"); diff --git a/src/main/java/org/rumbledb/spark/ml/ApplyEstimatorRuntimeIterator.java b/src/main/java/org/rumbledb/spark/ml/ApplyEstimatorRuntimeIterator.java index aa4b4aef60..f50bdb07f3 100644 --- a/src/main/java/org/rumbledb/spark/ml/ApplyEstimatorRuntimeIterator.java +++ b/src/main/java/org/rumbledb/spark/ml/ApplyEstimatorRuntimeIterator.java @@ -1,5 +1,6 @@ package org.rumbledb.spark.ml; +import lombok.extern.log4j.Log4j2; import org.apache.commons.lang3.StringUtils; import org.apache.spark.ml.Estimator; import org.apache.spark.ml.Transformer; @@ -36,6 +37,7 @@ import java.util.regex.Pattern; +@Log4j2 public class ApplyEstimatorRuntimeIterator extends AtMostOneItemLocalRuntimeIterator { private static final long serialVersionUID = 1L; @@ -81,8 +83,7 @@ public Item materializeFirstItemOrNull( } catch (IllegalArgumentException | NoSuchElementException e) { String message = e.getMessage(); if (message == null) { - System.err.println("Exception stack trace:"); - e.printStackTrace(); + log.error("Estimator fit failed with no exception message.", e); RumbleException ex = new InvalidRumbleMLParamException( "Parameters provided to " + this.estimatorShortName diff --git a/src/main/java/org/rumbledb/types/ItemTypeFactory.java b/src/main/java/org/rumbledb/types/ItemTypeFactory.java index 94d54c41c9..9d722b7b85 100644 --- a/src/main/java/org/rumbledb/types/ItemTypeFactory.java +++ b/src/main/java/org/rumbledb/types/ItemTypeFactory.java @@ -7,7 +7,7 @@ import java.util.List; import java.util.Map; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.ml.linalg.VectorUDT; import org.apache.spark.sql.types.ArrayType; import org.apache.spark.sql.types.CharType; @@ -28,6 +28,8 @@ import org.rumbledb.spark.SparkSessionManager; import org.rumbledb.runtime.typing.TypeInferrenceUtils; + +@Log4j2 public class ItemTypeFactory { public static ItemType createItemTypeFromJSoundCompactItem(Name name, Item item, StaticContext staticContext) { @@ -242,10 +244,9 @@ public static ItemType createItemTypeFromJSoundVerboseItem(Name name, Item item, Item contentItem = item.getItemByKey("content"); if (!keys.contains("content")) { - LogManager.getLogger("ItemTypeFactory") - .warn( - "The content facet of an object type is missing. By default, no fields are defined or overriden." - ); + log.warn( + "The content facet of an object type is missing. By default, no fields are defined or overriden." + ); contentItem = ItemFactory.getInstance().createArrayItem(staticContext.isQuerySideEffecting()); } else { if (contentItem == null) { @@ -270,10 +271,9 @@ public static ItemType createItemTypeFromJSoundVerboseItem(Name name, Item item, if (closedItem != null) { closed = closedItem.getBooleanValue(); } else { - LogManager.getLogger("ItemTypeFactory") - .warn( - "The closed facet of an object type is missing. By default, a closed object type is created. Set closed to false to keep the type open and allow arbitrary fields." - ); + log.warn( + "The closed facet of an object type is missing. By default, a closed object type is created. Set closed to false to keep the type open and allow arbitrary fields." + ); closed = true; } List contents = contentItem.getItemMembers(); diff --git a/src/main/java/org/rumbledb/types/SequenceType.java b/src/main/java/org/rumbledb/types/SequenceType.java index 7b8e770158..b1bcf6c96a 100644 --- a/src/main/java/org/rumbledb/types/SequenceType.java +++ b/src/main/java/org/rumbledb/types/SequenceType.java @@ -20,7 +20,7 @@ package org.rumbledb.types; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.rumbledb.context.DynamicContext; import org.rumbledb.context.Name; import org.rumbledb.context.StaticContext; @@ -34,6 +34,8 @@ import java.util.HashMap; import java.util.Map; + +@Log4j2 public class SequenceType implements Serializable { private static final long serialVersionUID = 1L; @@ -49,21 +51,19 @@ public SequenceType(ItemType itemType, Arity arity) { this.itemType = itemType; this.arity = arity; if (this.itemType == null) { - LogManager.getLogger("SequenceType") - .warn( - "Missing item type in incomplete sequence type " - + this.arity - + ", defaulting to item. Please let us know as we would like to look into this!" - ); + log.warn( + "Missing item type in incomplete sequence type " + + this.arity + + ", defaulting to item. Please let us know as we would like to look into this!" + ); this.itemType = BuiltinTypesCatalogue.item; } if (this.arity == null) { - LogManager.getLogger("SequenceType") - .warn( - "Missing arity in incomplete sequence type " - + this.itemType - + ", defaulting to *. Please let us know as we would like to look into this!" - ); + log.warn( + "Missing arity in incomplete sequence type " + + this.itemType + + ", defaulting to *. Please let us know as we would like to look into this!" + ); this.arity = Arity.ZeroOrMore; } } diff --git a/src/test/java/iq/SequentialClassificationTests.java b/src/test/java/iq/SequentialClassificationTests.java index e5f3101f5e..9519de154b 100644 --- a/src/test/java/iq/SequentialClassificationTests.java +++ b/src/test/java/iq/SequentialClassificationTests.java @@ -2,6 +2,7 @@ import iq.base.TestConfigurations; import iq.base.TestFileDiscovery; +import lombok.extern.log4j.Log4j2; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; @@ -36,6 +37,7 @@ import java.io.IOException; import java.net.URI; +@Log4j2 public class SequentialClassificationTests { private static final RumbleConfiguration configuration = TestConfigurations.defaultConfiguration(); @@ -309,7 +311,7 @@ public void testNonSequential() throws Throwable { ); int testIndex = 0; for (File testFile : TestFileDiscovery.jsoniqFiles(nonsequentialTestsDirectory)) { - System.err.println(testIndex++ + " : " + testFile); + log.debug("{} : {}", testIndex++, testFile); MainModule mainModule = parseAndCompile(testFile.getAbsolutePath()); for (Node descendant : mainModule.getDescendants()) { if (descendant instanceof Expression) { diff --git a/src/test/java/iq/UpdatesForRumbleBenchmark.java b/src/test/java/iq/UpdatesForRumbleBenchmark.java index 843a1c2e2a..6703792ff7 100644 --- a/src/test/java/iq/UpdatesForRumbleBenchmark.java +++ b/src/test/java/iq/UpdatesForRumbleBenchmark.java @@ -1,7 +1,7 @@ package iq; import org.apache.commons.io.FileUtils; -import org.apache.log4j.LogManager; +import lombok.extern.log4j.Log4j2; import org.apache.spark.SparkConf; import org.junit.jupiter.api.Assertions; import org.rumbledb.api.Item; @@ -24,6 +24,8 @@ import java.util.*; import java.util.function.Consumer; + +@Log4j2 public class UpdatesForRumbleBenchmark { private static final String APP_NAME = "Rumble application"; @@ -340,15 +342,14 @@ public UpdatesForRumbleBenchmark() { public static void setupSparkSession() { SparkSessionManager.getInstance().resetSession(); - System.err.println("Java version: " + javaVersion); - System.err.println("Scala version: " + scalaVersion); + log.info("Java version: {}", javaVersion); + log.info("Scala version: {}", scalaVersion); SparkConf sparkConfiguration = new SparkConf(); if (sparkConfiguration.get("spark.app.name", "").equals(" benchmarkDeltaTest(Rumble rumble, URI uri) throws IOException { @@ -567,9 +568,9 @@ public void deleteTable(String path) { try { File oldTable = new File(tableURI.getPath()); FileUtils.deleteDirectory(oldTable); - System.err.println("Deleted file: " + oldTable.getAbsolutePath()); + log.info("Deleted file: {}", oldTable.getAbsolutePath()); } catch (IOException e) { - e.printStackTrace(); + log.error("Could not delete old table.", e); Assertions.fail(); } } @@ -619,9 +620,9 @@ public static void main(String[] args) throws IOException { } } - System.out.println("##########################################"); - System.out.println("DONE"); - System.out.println("##########################################"); + log.info("##########################################"); + log.info("DONE"); + log.info("##########################################"); } diff --git a/src/test/java/iq/base/AnnotationTestExecutor.java b/src/test/java/iq/base/AnnotationTestExecutor.java index 29f423d5e4..8d84fd5ff1 100644 --- a/src/test/java/iq/base/AnnotationTestExecutor.java +++ b/src/test/java/iq/base/AnnotationTestExecutor.java @@ -17,6 +17,7 @@ package iq.base; +import lombok.extern.log4j.Log4j2; import org.apache.commons.lang3.exception.ExceptionUtils; import org.junit.jupiter.api.Assertions; import org.rumbledb.api.ExternalBindings; @@ -40,6 +41,7 @@ import java.io.Reader; import java.net.URI; +@Log4j2 public final class AnnotationTestExecutor { private AnnotationTestExecutor() { @@ -170,7 +172,6 @@ private static void assertExpectedFailure( annotation.errorCode(), annotation.errorMetadata() ); - System.out.println(executionResult.failureMessage()); } private static String unexpectedFailureMessage( @@ -302,7 +303,7 @@ private static String formatSequenceForLegacyRuntimeAssertions( sb.append(")"); if (sequence.hasNext() && resultSizeCap > 0 && itemCount == resultSizeCap) { - System.err.println( + log.warn( "Warning! The output sequence contains a large number of items but its materialization was capped at " + resultSizeCap + " items. This value can be configured with the --result-size parameter at startup" diff --git a/src/test/java/iq/base/SparkAnnotationsTestsBase.java b/src/test/java/iq/base/SparkAnnotationsTestsBase.java index 136b72ce93..558bbf1a44 100644 --- a/src/test/java/iq/base/SparkAnnotationsTestsBase.java +++ b/src/test/java/iq/base/SparkAnnotationsTestsBase.java @@ -17,18 +17,20 @@ package iq.base; +import lombok.extern.log4j.Log4j2; import org.apache.spark.SparkConf; import org.junit.jupiter.api.BeforeAll; import scala.util.Properties; import org.rumbledb.spark.SparkSessionManager; +@Log4j2 public abstract class SparkAnnotationsTestsBase extends AnnotationsTestsBase { @BeforeAll final void setupSparkSession() { SparkSessionManager.getInstance().resetSession(); - System.err.println("Java version: " + System.getProperty("java.version")); - System.err.println("Scala version: " + Properties.scalaPropOrElse("version.number", () -> "unknown")); + log.info("Java version: {}", System.getProperty("java.version")); + log.info("Scala version: {}", Properties.scalaPropOrElse("version.number", () -> "unknown")); SparkConf sparkConfiguration = new SparkConf(); sparkConfiguration.setMaster("local[*]"); @@ -40,7 +42,7 @@ final void setupSparkSession() { configureSpark(sparkConfiguration); SparkSessionManager.getInstance().initializeConfigurationAndSession(sparkConfiguration, true); - System.err.println("Spark version: " + SparkSessionManager.getInstance().getJavaSparkContext().version()); + log.info("Spark version: {}", SparkSessionManager.getInstance().getJavaSparkContext().version()); } protected void configureSpark(SparkConf sparkConfiguration) {