Эх сурвалжийг харах

[spark] Rename cell to statement

Erick Tryzelaar 11 жил өмнө
parent
commit
6304ca8

+ 2 - 2
apps/spark/java/sparker-server/src/main/java/com/cloudera/hue/sparker/server/SparkerApp.java

@@ -1,6 +1,6 @@
 package com.cloudera.hue.sparker.server;
 
-import com.cloudera.hue.sparker.server.resources.CellResource;
+import com.cloudera.hue.sparker.server.resources.StatementResource;
 import com.cloudera.hue.sparker.server.resources.SessionResource;
 import com.cloudera.hue.sparker.server.sessions.SessionManager;
 import com.sun.jersey.core.spi.factory.ResponseBuilderImpl;
@@ -26,7 +26,7 @@ public class SparkerApp extends Application<SparkerConfiguration> {
     public void run(SparkerConfiguration sparkerConfiguration, Environment environment) throws Exception {
         final SessionManager sessionManager = new SessionManager();
         environment.jersey().register(new SessionResource(sessionManager));
-        environment.jersey().register(new CellResource(sessionManager));
+        environment.jersey().register(new StatementResource(sessionManager));
         environment.jersey().register(new SessionManagerExceptionMapper());
     }
 

+ 0 - 53
apps/spark/java/sparker-server/src/main/java/com/cloudera/hue/sparker/server/resources/CellResource.java

@@ -1,53 +0,0 @@
-package com.cloudera.hue.sparker.server.resources;
-
-import com.cloudera.hue.sparker.server.sessions.Cell;
-import com.cloudera.hue.sparker.server.sessions.Session;
-import com.cloudera.hue.sparker.server.sessions.SessionManager;
-import com.codahale.metrics.annotation.Timed;
-
-import javax.ws.rs.*;
-import javax.ws.rs.core.MediaType;
-import java.util.List;
-
-@Path("/sessions/{sessionId}/cells")
-@Produces(MediaType.APPLICATION_JSON)
-public class CellResource {
-
-    private final SessionManager sessionManager;
-
-    public CellResource(SessionManager sessionManager) {
-        this.sessionManager = sessionManager;
-    }
-
-    @GET
-    @Timed
-    public List<Cell> getCells(@PathParam("sessionId") String sessionId,
-                               @QueryParam("from") Integer fromCell,
-                               @QueryParam("limit") Integer limit) throws SessionManager.SessionNotFound {
-        Session session = sessionManager.get(sessionId);
-        List<Cell> cells = session.getCells();
-
-        if (fromCell != null || limit != null) {
-            if (fromCell == null) {
-                fromCell = 0;
-            }
-
-            if (limit == null) {
-                limit = cells.size();
-            }
-
-            cells = cells.subList(fromCell, fromCell + limit);
-        }
-
-        return cells;
-    }
-
-    @Path("/{cellId}")
-    @GET
-    @Timed
-    public Cell getCell(@PathParam("sessionId") String sessionId, @PathParam("cellId") int cellId) throws SessionManager.SessionNotFound {
-        Session session = sessionManager.get(sessionId);
-        return session.getCell(cellId);
-    }
-
-}

+ 8 - 9
apps/spark/java/sparker-server/src/main/java/com/cloudera/hue/sparker/server/resources/SessionResource.java

@@ -1,12 +1,11 @@
 package com.cloudera.hue.sparker.server.resources;
 
-import com.cloudera.hue.sparker.server.sessions.Cell;
+import com.cloudera.hue.sparker.server.sessions.Statement;
 import com.cloudera.hue.sparker.server.sessions.ClosedSessionException;
 import com.cloudera.hue.sparker.server.sessions.Session;
 import com.cloudera.hue.sparker.server.sessions.SessionManager;
 import com.codahale.metrics.annotation.Timed;
 import com.sun.jersey.core.spi.factory.ResponseBuilderImpl;
-import org.hibernate.validator.constraints.NotEmpty;
 
 import javax.servlet.http.HttpServletRequest;
 import javax.validation.Valid;
@@ -76,11 +75,11 @@ public class SessionResource {
     @Path("/{id}")
     @GET
     @Timed
-    public List<Cell> getSession(@PathParam("id") String id,
+    public List<Statement> getSession(@PathParam("id") String id,
                                  @QueryParam("from") Integer fromCell,
                                  @QueryParam("limit") Integer limit) throws SessionManager.SessionNotFound {
         Session session = sessionManager.get(id);
-        List<Cell> cells = session.getCells();
+        List<Statement> statements = session.getStatements();
 
         if (fromCell != null || limit != null) {
             if (fromCell == null) {
@@ -88,13 +87,13 @@ public class SessionResource {
             }
 
             if (limit == null) {
-                limit = cells.size();
+                limit = statements.size();
             }
 
-            cells = cells.subList(fromCell, fromCell + limit);
+            statements = statements.subList(fromCell, fromCell + limit);
         }
 
-        return cells;
+        return statements;
     }
 
     @Path("/{id}")
@@ -106,9 +105,9 @@ public class SessionResource {
         Session session = sessionManager.get(id);
 
         // The cell is evaluated inline, but eventually it'll be turned into an asynchronous call.
-        Cell cell = session.executeStatement(body.getStatement());
+        Statement statement = session.executeStatement(body.getStatement());
 
-        URI location = new URI("/cells/" + cell.getId());
+        URI location = new URI("/cells/" + statement.getId());
         return Response.created(location).build();
     }
 

+ 47 - 0
apps/spark/java/sparker-server/src/main/java/com/cloudera/hue/sparker/server/resources/StatementResource.java

@@ -0,0 +1,47 @@
+package com.cloudera.hue.sparker.server.resources;
+
+import com.cloudera.hue.sparker.server.sessions.Session;
+import com.cloudera.hue.sparker.server.sessions.SessionManager;
+import com.cloudera.hue.sparker.server.sessions.Statement;
+import com.codahale.metrics.annotation.Timed;
+
+import javax.ws.rs.*;
+import javax.ws.rs.core.MediaType;
+import java.util.List;
+
+@Path("/sessions/{sessionId}/statements")
+@Produces(MediaType.APPLICATION_JSON)
+public class StatementResource {
+
+    private final SessionManager sessionManager;
+
+    public StatementResource(SessionManager sessionManager) {
+        this.sessionManager = sessionManager;
+    }
+
+    @GET
+    @Timed
+    public List<Statement> getStatements(@PathParam("sessionId") String sessionId,
+                                         @QueryParam("from") Integer fromStatement,
+                                         @QueryParam("limit") Integer limit) throws SessionManager.SessionNotFound {
+        Session session = sessionManager.get(sessionId);
+        List<Statement> statements;
+
+        if (fromStatement == null && limit == null) {
+            statements = session.getStatements();
+        } else {
+            statements = session.getStatementRange(fromStatement, fromStatement + limit);
+        }
+
+        return statements;
+    }
+
+    @Path("/{statementId}")
+    @GET
+    @Timed
+    public Statement getStatement(@PathParam("sessionId") String sessionId, @PathParam("statementId") int statementId) throws SessionManager.SessionNotFound, Session.StatementNotFound {
+        Session session = sessionManager.get(sessionId);
+        return session.getStatement(statementId);
+    }
+
+}

+ 8 - 4
apps/spark/java/sparker-server/src/main/java/com/cloudera/hue/sparker/server/sessions/Session.java

@@ -30,17 +30,21 @@ public interface Session {
     String getId();
 
     @JsonProperty
-    List<Cell> getCells();
+    List<Statement> getStatements();
 
-    List<Cell> getCellRange(int fromIndex, int toIndex);
+    List<Statement> getStatementRange(Integer fromIndex, Integer toIndex);
 
-    Cell getCell(int cellId);
+    Statement getStatement(int statementId) throws StatementNotFound;
 
     @JsonProperty
     public long getLastActivity();
 
-    public Cell executeStatement(String statement) throws Exception, ClosedSessionException;
+    public Statement executeStatement(String statement) throws Exception, ClosedSessionException;
 
     public void close() throws IOException, InterruptedException, TimeoutException;
+
+    public static class StatementNotFound extends Throwable {
+
+    }
 }
 

+ 15 - 15
apps/spark/java/sparker-server/src/main/java/com/cloudera/hue/sparker/server/sessions/SparkSession.java

@@ -58,7 +58,7 @@ public class SparkSession implements Session {
     private final Process process;
     private final Writer writer;
     private final BufferedReader reader;
-    private final List<Cell> cells = new ArrayList<Cell>();
+    private final List<Statement> statements = new ArrayList<Statement>();
     private final ObjectMapper objectMapper = new ObjectMapper();
 
     private boolean isClosed = false;
@@ -93,36 +93,36 @@ public class SparkSession implements Session {
     }
 
     @Override
-    synchronized public List<Cell> getCells() {
-        return Lists.newArrayList(cells);
+    synchronized public List<Statement> getStatements() {
+        return Lists.newArrayList(statements);
     }
 
     @Override
-    synchronized public List<Cell> getCellRange(int fromIndex, int toIndex) {
-        return cells.subList(fromIndex, toIndex);
+    synchronized public List<Statement> getStatementRange(Integer fromIndex, Integer toIndex) {
+        return statements.subList(fromIndex, toIndex);
     }
 
     @Override
-    synchronized public Cell getCell(int index) {
-        return cells.get(index);
+    synchronized public Statement getStatement(int index) {
+        return statements.get(index);
     }
 
     @Override
-    synchronized public Cell executeStatement(String statement) throws IOException, ClosedSessionException, InterruptedException {
+    synchronized public Statement executeStatement(String statementStr) throws IOException, ClosedSessionException, InterruptedException {
         if (isClosed) {
             throw new ClosedSessionException();
         }
 
         touchLastActivity();
 
-        Cell cell = new Cell(cells.size());
-        cells.add(cell);
+        Statement statement = new Statement(statements.size());
+        statements.add(statement);
 
-        cell.addInput(statement);
+        statement.addInput(statementStr);
 
         ObjectNode request = objectMapper.createObjectNode();
         request.put("type", "stdin");
-        request.put("statement", statement);
+        request.put("statement", statementStr);
 
         writer.write(request.toString());
         writer.write("\n");
@@ -141,14 +141,14 @@ public class SparkSession implements Session {
         JsonNode response = objectMapper.readTree(line);
 
         if (response.has("stdout")) {
-            cell.addOutput(response.get("stdout").asText());
+            statement.addOutput(response.get("stdout").asText());
         }
 
         if (response.has("stderr")) {
-            cell.addOutput(response.get("stderr").asText());
+            statement.addOutput(response.get("stderr").asText());
         }
 
-        return cell;
+        return statement;
     }
 
     @Override

+ 2 - 2
apps/spark/java/sparker-server/src/main/java/com/cloudera/hue/sparker/server/sessions/Cell.java → apps/spark/java/sparker-server/src/main/java/com/cloudera/hue/sparker/server/sessions/Statement.java

@@ -5,7 +5,7 @@ import com.fasterxml.jackson.annotation.JsonProperty;
 import java.util.ArrayList;
 import java.util.List;
 
-public class Cell {
+public class Statement {
 
     public enum State {
         NOT_READY,
@@ -22,7 +22,7 @@ public class Cell {
 
     final List<String> error = new ArrayList<String>();
 
-    public Cell(int id) {
+    public Statement(int id) {
         this.id = id;
         this.state = State.COMPLETE;
     }