Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 13 additions & 15 deletions .github/workflows/core.yml
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ jobs:
strategy:
fail-fast: false
matrix:
java: [ 11 ]
java: [ 11, 17 ]
steps:
- name: Checkout
uses: actions/checkout@v5
Expand Down Expand Up @@ -90,7 +90,7 @@ jobs:
strategy:
fail-fast: false
matrix:
java: [ 11 ]
java: [ 11, 17 ]
env:
INTERPRETERS: 'hbase,jdbc,file,flink-cmd,cassandra,elasticsearch,bigquery,livy,groovy,java,neo4j,sparql,mongodb,influxdb,shell'
steps:
Expand Down Expand Up @@ -136,7 +136,7 @@ jobs:
fail-fast: false
matrix:
python: [ 3.9 ]
java: [ 11 ]
java: [ 11, 17 ]
steps:
- name: Checkout
uses: actions/checkout@v5
Expand Down Expand Up @@ -181,7 +181,7 @@ jobs:
strategy:
fail-fast: false
matrix:
java: [ 11 ]
java: [ 11, 17 ]
steps:
# user/password => root/root
- name: Start mysql
Expand Down Expand Up @@ -232,22 +232,24 @@ jobs:
fail-fast: false
matrix:
include:
- python: 3.9
- java: 11
python: 3.9
flink: 119
flink-profile: "1.19"
- python: 3.9
- java: 17
python: 3.9
flink: 120
flink-profile: "1.20"
steps:
- name: Checkout
uses: actions/checkout@v5
- name: Tune Runner VM
uses: ./.github/actions/tune-runner-vm
- name: Set up JDK 11
- name: Set up JDK ${{ matrix.java }}
uses: actions/setup-java@v5
with:
distribution: 'temurin'
java-version: 11
java-version: ${{ matrix.java }}
- name: Cache local Maven repository
uses: actions/cache@v5
with:
Expand Down Expand Up @@ -285,7 +287,7 @@ jobs:
strategy:
fail-fast: false
matrix:
java: [ 11 ]
java: [ 11, 17 ]
steps:
- name: Checkout
uses: actions/checkout@v5
Expand Down Expand Up @@ -366,22 +368,18 @@ jobs:
auto-activate: false
use-mamba: true
- name: run spark-3.3 tests with scala-2.12 and python-${{ matrix.python }}
if: ${{ matrix.java == 11 }}
run: |
rm -rf spark/interpreter/metastore_db
./mvnw verify -pl spark-submit,spark/interpreter -am -Dtest=org/apache/zeppelin/spark/* -Pspark-3.3 -Pspark-scala-2.12 -Pintegration -DfailIfNoTests=false ${MAVEN_ARGS}
- name: run spark-3.3 tests with scala-2.13 and python-${{ matrix.python }}
if: ${{ matrix.java == 11 }}
run: |
rm -rf spark/interpreter/metastore_db
./mvnw verify -pl spark-submit,spark/interpreter -am -Dtest=org/apache/zeppelin/spark/* -Pspark-3.3 -Pspark-scala-2.13 -Pintegration -DfailIfNoTests=false ${MAVEN_ARGS}
- name: run spark-3.4 tests with scala-2.13 and python-${{ matrix.python }}
if: ${{ matrix.java == 11 }}
run: |
rm -rf spark/interpreter/metastore_db
./mvnw verify -pl spark-submit,spark/interpreter -am -Dtest=org/apache/zeppelin/spark/* -Pspark-3.4 -Pspark-scala-2.13 -Pintegration -DfailIfNoTests=false ${MAVEN_ARGS}
- name: run spark-3.5 tests with scala-2.13 and python-${{ matrix.python }}
if: ${{ matrix.java == 11 }}
run: |
rm -rf spark/interpreter/metastore_db
./mvnw verify -pl spark-submit,spark/interpreter -am -Dtest=org/apache/zeppelin/spark/* -Pspark-3.5 -Pspark-scala-2.13 -Pintegration -DfailIfNoTests=false ${MAVEN_ARGS}
Expand Down Expand Up @@ -442,7 +440,7 @@ jobs:
strategy:
fail-fast: false
matrix:
java: [ 11 ]
java: [ 11, 17 ]
steps:
- name: Checkout
uses: actions/checkout@v5
Expand Down Expand Up @@ -472,7 +470,7 @@ jobs:
strategy:
fail-fast: false
matrix:
java: [ 11 ]
java: [ 17 ]
steps:
- name: Checkout
uses: actions/checkout@v5
Expand Down
26 changes: 25 additions & 1 deletion bin/common.cmd
Original file line number Diff line number Diff line change
Expand Up @@ -71,14 +71,38 @@ if not defined ZEPPELIN_JAVA_OPTS (
set ZEPPELIN_JAVA_OPTS=%ZEPPELIN_JAVA_OPTS% -Dfile.encoding=%ZEPPELIN_ENCODING% %ZEPPELIN_MEM%
)

REM JPMS (Java Platform Module System) args mirrored from pom.xml extraJavaTestArgs.
REM Targets JDK 17+: --add-modules, --enable-native-access, --sun-misc-unsafe-memory-access
REM require JDK 17+. -XX:+IgnoreUnrecognizedVMOptions only silences unknown -XX flags.
set JPMS_JAVA_OPTS=-XX:+IgnoreUnrecognizedVMOptions
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% --add-opens=java.base/java.lang=ALL-UNNAMED
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% --add-opens=java.base/java.lang.invoke=ALL-UNNAMED
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% --add-opens=java.base/java.lang.reflect=ALL-UNNAMED
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% --add-opens=java.base/java.io=ALL-UNNAMED
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% --add-opens=java.base/java.net=ALL-UNNAMED
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% --add-opens=java.base/java.nio=ALL-UNNAMED
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% --add-opens=java.base/java.util=ALL-UNNAMED
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% --add-opens=java.base/java.util.concurrent=ALL-UNNAMED
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% --add-opens=java.base/java.util.concurrent.atomic=ALL-UNNAMED
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% --add-opens=java.base/jdk.internal.ref=ALL-UNNAMED
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% --add-opens=java.base/sun.nio.ch=ALL-UNNAMED
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% --add-opens=java.base/sun.nio.cs=ALL-UNNAMED
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% --add-opens=java.base/sun.security.action=ALL-UNNAMED
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% --add-opens=java.base/sun.util.calendar=ALL-UNNAMED
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% -Dio.netty.tryReflectionSetAccessible=true
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% -Dio.netty.allocator.type=pooled
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% -Dio.netty.handler.ssl.defaultEndpointVerificationAlgorithm=NONE
set JPMS_JAVA_OPTS=%JPMS_JAVA_OPTS% --sun-misc-unsafe-memory-access=allow --enable-native-access=ALL-UNNAMED

if not defined JAVA_OPTS (
set JAVA_OPTS=%ZEPPELIN_JAVA_OPTS%
) else (
set JAVA_OPTS=%JAVA_OPTS% %ZEPPELIN_JAVA_OPTS%
)
set JAVA_OPTS=%JAVA_OPTS% %JPMS_JAVA_OPTS%


set JAVA_INTP_OPTS=%ZEPPELIN_INTP_JAVA_OPTS% -Dfile.encoding=%ZEPPELIN_ENCODING%
set JAVA_INTP_OPTS=%ZEPPELIN_INTP_JAVA_OPTS% -Dfile.encoding=%ZEPPELIN_ENCODING% %JPMS_JAVA_OPTS%

if not defined JAVA_HOME (
set ZEPPELIN_RUNNER=java
Expand Down
26 changes: 24 additions & 2 deletions bin/common.sh
Original file line number Diff line number Diff line change
Expand Up @@ -147,15 +147,37 @@ if [[ ( -z "${ZEPPELIN_INTP_MEM}" ) && ( "${ZEPPELIN_INTERPRETER_LAUNCHER}" != "
export ZEPPELIN_INTP_MEM="-Xmx1024m"
fi

JAVA_OPTS+=" ${ZEPPELIN_JAVA_OPTS} -Dfile.encoding=${ZEPPELIN_ENCODING} ${ZEPPELIN_MEM}"
# JPMS (Java Platform Module System) args mirrored from pom.xml extraJavaTestArgs.
# Targets JDK 17+: --add-modules, --enable-native-access, --sun-misc-unsafe-memory-access
# require JDK 17+. -XX:+IgnoreUnrecognizedVMOptions only silences unknown -XX flags.
JPMS_JAVA_OPTS="-XX:+IgnoreUnrecognizedVMOptions"
JPMS_JAVA_OPTS+=" --add-opens=java.base/java.lang=ALL-UNNAMED"
JPMS_JAVA_OPTS+=" --add-opens=java.base/java.lang.invoke=ALL-UNNAMED"
JPMS_JAVA_OPTS+=" --add-opens=java.base/java.lang.reflect=ALL-UNNAMED"
JPMS_JAVA_OPTS+=" --add-opens=java.base/java.io=ALL-UNNAMED"
JPMS_JAVA_OPTS+=" --add-opens=java.base/java.net=ALL-UNNAMED"
JPMS_JAVA_OPTS+=" --add-opens=java.base/java.nio=ALL-UNNAMED"
JPMS_JAVA_OPTS+=" --add-opens=java.base/java.util=ALL-UNNAMED"
JPMS_JAVA_OPTS+=" --add-opens=java.base/java.util.concurrent=ALL-UNNAMED"
JPMS_JAVA_OPTS+=" --add-opens=java.base/java.util.concurrent.atomic=ALL-UNNAMED"
JPMS_JAVA_OPTS+=" --add-opens=java.base/jdk.internal.ref=ALL-UNNAMED"
JPMS_JAVA_OPTS+=" --add-opens=java.base/sun.nio.ch=ALL-UNNAMED"
JPMS_JAVA_OPTS+=" --add-opens=java.base/sun.nio.cs=ALL-UNNAMED"
JPMS_JAVA_OPTS+=" --add-opens=java.base/sun.security.action=ALL-UNNAMED"
JPMS_JAVA_OPTS+=" --add-opens=java.base/sun.util.calendar=ALL-UNNAMED"
JPMS_JAVA_OPTS+=" -Dio.netty.tryReflectionSetAccessible=true"
JPMS_JAVA_OPTS+=" -Dio.netty.allocator.type=pooled"
JPMS_JAVA_OPTS+=" -Dio.netty.handler.ssl.defaultEndpointVerificationAlgorithm=NONE"
JPMS_JAVA_OPTS+=" --sun-misc-unsafe-memory-access=allow --enable-native-access=ALL-UNNAMED"
JAVA_OPTS+=" ${ZEPPELIN_JAVA_OPTS} -Dfile.encoding=${ZEPPELIN_ENCODING} ${ZEPPELIN_MEM} ${JPMS_JAVA_OPTS}"
if [[ -n "${ZEPPELIN_IN_DOCKER}" ]]; then
JAVA_OPTS+=" -Dlog4j.configuration=file://${ZEPPELIN_CONF_DIR}/log4j_docker.properties"
else
JAVA_OPTS+=" -Dlog4j.configuration=file://${ZEPPELIN_CONF_DIR}/log4j.properties"
fi
export JAVA_OPTS

JAVA_INTP_OPTS="${ZEPPELIN_INTP_JAVA_OPTS} -Dfile.encoding=${ZEPPELIN_ENCODING}"
JAVA_INTP_OPTS="${ZEPPELIN_INTP_JAVA_OPTS} -Dfile.encoding=${ZEPPELIN_ENCODING} ${JPMS_JAVA_OPTS}"
if [[ -n "${ZEPPELIN_IN_DOCKER}" ]]; then
JAVA_INTP_OPTS+=" -Dlog4j.configuration=file://${ZEPPELIN_CONF_DIR}/log4j_docker.properties -Dlog4j.configurationFile=file://${ZEPPELIN_CONF_DIR}/log4j2_docker.properties"
elif [[ -z "${ZEPPELIN_SPARK_YARN_CLUSTER}" ]]; then
Expand Down
21 changes: 5 additions & 16 deletions cassandra/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -37,8 +37,6 @@
<scalate.version>1.9.8</scalate.version>

<!-- test library versions -->
<jna.version>5.12.1</jna.version>
<cassandra.unit.version>4.3.1.0</cassandra.unit.version>

<scala.version>${scala.2.12.version}</scala.version>
<scala.binary.version>2.12</scala.binary.version>
Expand Down Expand Up @@ -137,26 +135,17 @@
</dependency>

<dependency>
<groupId>net.java.dev.jna</groupId>
<artifactId>jna</artifactId>
<version>${jna.version}</version>
<groupId>org.testcontainers</groupId>
<artifactId>cassandra</artifactId>
<scope>test</scope>
</dependency>

<dependency>
<groupId>org.cassandraunit</groupId>
<artifactId>cassandra-unit</artifactId>
<version>${cassandra.unit.version}</version>
<scope>test</scope>
<exclusions>
<exclusion>
<groupId>com.datastax.oss</groupId>
<artifactId>java-driver-core</artifactId>
</exclusion>
</exclusions>
<groupId>org.testcontainers</groupId>
<artifactId>junit-jupiter</artifactId>
<scope>test</scope>
</dependency>


<dependency>
<groupId>org.mockito</groupId>
<artifactId>mockito-core</artifactId>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,13 +27,13 @@
import org.apache.zeppelin.interpreter.InterpreterContext;
import org.apache.zeppelin.interpreter.InterpreterResult;
import org.apache.zeppelin.interpreter.InterpreterResult.Code;
import org.cassandraunit.CQLDataLoader;
import org.cassandraunit.dataset.cql.ClassPathCQLDataSet;
import org.cassandraunit.utils.EmbeddedCassandraServerHelper;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Disabled;
import org.junit.jupiter.api.Test;
import org.testcontainers.containers.CassandraContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;

import java.io.IOException;
import java.nio.charset.StandardCharsets;
Expand Down Expand Up @@ -65,28 +65,40 @@
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;

public class CassandraInterpreterTest { // extends AbstractCassandraUnit4CQLTestCase {
@Testcontainers
public class CassandraInterpreterTest {
private static final String ARTISTS_TABLE = "zeppelin.artists";

private static volatile CassandraInterpreter interpreter;

private static CqlSession session;

private final InterpreterContext intrContext = InterpreterContext.builder()
.setParagraphTitle("Paragraph1")
.build();

@Container
public static CassandraContainer<?> cassandra =
new CassandraContainer<>("cassandra:4.1.3");

@BeforeAll
public static synchronized void setUp() throws IOException, InterruptedException {
System.setProperty("cassandra.skip_wait_for_gossip_to_settle", "0");
System.setProperty("cassandra.load_ring_state", "false");
System.setProperty("cassandra.initial_token", "0");
System.setProperty("cassandra.num_tokens", "nil");
System.setProperty("cassandra.allocate_tokens_for_local_replication_factor", "nil");
EmbeddedCassandraServerHelper.startEmbeddedCassandra();
CqlSession session = EmbeddedCassandraServerHelper.getSession();
new CQLDataLoader(session).load(new ClassPathCQLDataSet("prepare_all.cql", "zeppelin"));
public static synchronized void setUp() throws IOException {
session = CqlSession.builder()
.addContactPoint(java.net.InetSocketAddress.createUnresolved(
cassandra.getHost(), cassandra.getMappedPort(9042)))
.withLocalDatacenter("datacenter1")
.build();

String cql = IOUtils.resourceToString("/prepare_all.cql", StandardCharsets.UTF_8);
for (String stmt : cql.split(";")) {
String trimmed = stmt.trim();
if (!trimmed.isEmpty()) {
session.execute(trimmed);
}
}

Properties properties = new Properties();
properties.setProperty(CASSANDRA_CLUSTER_NAME, EmbeddedCassandraServerHelper.getClusterName());
properties.setProperty(CASSANDRA_CLUSTER_NAME, "Test Cluster");
properties.setProperty(CASSANDRA_COMPRESSION_PROTOCOL, "NONE");
properties.setProperty(CASSANDRA_CREDENTIALS_USERNAME, "none");
properties.setProperty(CASSANDRA_CREDENTIALS_PASSWORD, "none");
Expand All @@ -111,9 +123,9 @@ public static synchronized void setUp() throws IOException, InterruptedException
properties.setProperty(CASSANDRA_SOCKET_READ_TIMEOUT_MILLIS, "12000");
properties.setProperty(CASSANDRA_SOCKET_TCP_NO_DELAY, "true");

properties.setProperty(CASSANDRA_HOSTS, EmbeddedCassandraServerHelper.getHost());
properties.setProperty(CASSANDRA_HOSTS, cassandra.getHost());
properties.setProperty(CASSANDRA_PORT,
Integer.toString(EmbeddedCassandraServerHelper.getNativeTransportPort()));
Integer.toString(cassandra.getMappedPort(9042)));
properties.setProperty("datastax-java-driver.advanced.connection.pool.local.size", "1");
interpreter = new CassandraInterpreter(properties);
interpreter.open();
Expand All @@ -122,6 +134,9 @@ public static synchronized void setUp() throws IOException, InterruptedException
@AfterAll
public static void tearDown() {
interpreter.close();
if (session != null) {
session.close();
}
}

@Test
Expand Down Expand Up @@ -333,7 +348,7 @@ void should_execute_statement_with_timestamp_option() throws Exception {
String statement2 = "@timestamp=15\n" +
"INSERT INTO zeppelin.ts(key,val) VALUES('k','v2');";

CqlSession session = EmbeddedCassandraServerHelper.getSession();
CqlSession session = CassandraInterpreterTest.session;
// Insert v1 with current timestamp
interpreter.interpret(statement1, intrContext);
System.out.println("going to read data from zeppelin.ts;");
Expand Down Expand Up @@ -562,14 +577,17 @@ void should_display_statistics_for_non_select_statement() {

// When
final InterpreterResult actual = interpreter.interpret(query, intrContext);
final int port = EmbeddedCassandraServerHelper.getNativeTransportPort();
final String address = EmbeddedCassandraServerHelper.getHost();
final int port = cassandra.getMappedPort(9042);
final String address = cassandra.getHost();
// Then
final String expected = rawResult.replaceAll("TRIED_HOSTS", address + ":" + port)
.replaceAll("QUERIED_HOSTS", address + ":" + port);

assertEquals(Code.SUCCESS, actual.code());
assertEquals(expected, reformatHtml(actual.message().get(0).getData()));
// JDK 17+ renders unresolved InetSocketAddress as "host/<unresolved>:port"
String actualHtml = reformatHtml(actual.message().get(0).getData())
.replaceAll(address + "/&lt;unresolved&gt;:", address + ":");
assertEquals(expected, actualHtml);
}

@Test
Expand Down

Large diffs are not rendered by default.

Loading
Loading