Skip to content

Commit 567a261

Browse files
committed
add raw event API
Signed-off-by: Maciej Obuchowski <obuchowski.maciej@gmail.com>
1 parent b709b03 commit 567a261

24 files changed

Lines changed: 525 additions & 1 deletion

api/build.gradle

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@ dependencies {
4040
implementation "io.prometheus:simpleclient_hotspot:${prometheusVersion}"
4141
implementation "io.prometheus:simpleclient_servlet:${prometheusVersion}"
4242
implementation "org.jdbi:jdbi3-core:${jdbi3Version}"
43+
implementation "org.jdbi:jdbi3-jackson2:${jdbi3Version}"
4344
implementation "org.jdbi:jdbi3-postgres:${jdbi3Version}"
4445
implementation "org.jdbi:jdbi3-sqlobject:${jdbi3Version}"
4546
implementation 'com.google.guava:guava:31.1-jre'

api/src/main/java/marquez/MarquezContext.java

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@
1414
import lombok.Getter;
1515
import lombok.NonNull;
1616
import marquez.api.DatasetResource;
17+
import marquez.api.EventResource;
1718
import marquez.api.JobResource;
1819
import marquez.api.NamespaceResource;
1920
import marquez.api.OpenLineageResource;
@@ -25,6 +26,7 @@
2526
import marquez.db.DatasetDao;
2627
import marquez.db.DatasetFieldDao;
2728
import marquez.db.DatasetVersionDao;
29+
import marquez.db.EventDao;
2830
import marquez.db.JobContextDao;
2931
import marquez.db.JobDao;
3032
import marquez.db.JobVersionDao;
@@ -42,6 +44,7 @@
4244
import marquez.service.DatasetFieldService;
4345
import marquez.service.DatasetService;
4446
import marquez.service.DatasetVersionService;
47+
import marquez.service.EventService;
4548
import marquez.service.JobService;
4649
import marquez.service.LineageService;
4750
import marquez.service.NamespaceService;
@@ -71,6 +74,7 @@ public final class MarquezContext {
7174
@Getter private final OpenLineageDao openLineageDao;
7275
@Getter private final LineageDao lineageDao;
7376
@Getter private final SearchDao searchDao;
77+
@Getter private final EventDao eventDao;
7478

7579
@Getter private final List<RunTransitionListener> runTransitionListeners;
7680

@@ -82,6 +86,7 @@ public final class MarquezContext {
8286
@Getter private final RunService runService;
8387
@Getter private final OpenLineageService openLineageService;
8488
@Getter private final LineageService lineageService;
89+
@Getter private final EventService eventService;
8590

8691
@Getter private final NamespaceResource namespaceResource;
8792
@Getter private final SourceResource sourceResource;
@@ -90,6 +95,7 @@ public final class MarquezContext {
9095
@Getter private final TagResource tagResource;
9196
@Getter private final OpenLineageResource openLineageResource;
9297
@Getter private final SearchResource searchResource;
98+
@Getter private final EventResource eventResource;
9399

94100
@Getter private final ImmutableList<Object> resources;
95101
@Getter private final JdbiExceptionExceptionMapper jdbiException;
@@ -119,6 +125,7 @@ private MarquezContext(
119125
this.openLineageDao = jdbi.onDemand(OpenLineageDao.class);
120126
this.lineageDao = jdbi.onDemand(LineageDao.class);
121127
this.searchDao = jdbi.onDemand(SearchDao.class);
128+
this.eventDao = jdbi.onDemand(EventDao.class);
122129
this.runTransitionListeners = runTransitionListeners;
123130

124131
this.namespaceService = new NamespaceService(baseDao);
@@ -131,6 +138,7 @@ private MarquezContext(
131138
this.tagService.init(tags);
132139
this.openLineageService = new OpenLineageService(baseDao, runService);
133140
this.lineageService = new LineageService(lineageDao, jobDao);
141+
this.eventService = new EventService(baseDao);
134142
this.jdbiException = new JdbiExceptionExceptionMapper();
135143
final ServiceFactory serviceFactory =
136144
ServiceFactory.builder()
@@ -144,6 +152,7 @@ private MarquezContext(
144152
.lineageService(lineageService)
145153
.datasetFieldService(new DatasetFieldService(baseDao))
146154
.datasetVersionService(new DatasetVersionService(baseDao))
155+
.eventService(eventService)
147156
.build();
148157
this.namespaceResource = new NamespaceResource(serviceFactory);
149158
this.sourceResource = new SourceResource(serviceFactory);
@@ -152,6 +161,7 @@ private MarquezContext(
152161
this.tagResource = new TagResource(serviceFactory);
153162
this.openLineageResource = new OpenLineageResource(serviceFactory);
154163
this.searchResource = new SearchResource(searchDao);
164+
this.eventResource = new EventResource(serviceFactory);
155165

156166
this.resources =
157167
ImmutableList.of(
@@ -162,7 +172,8 @@ private MarquezContext(
162172
tagResource,
163173
jdbiException,
164174
openLineageResource,
165-
searchResource);
175+
searchResource,
176+
eventResource);
166177

167178
final MarquezGraphqlServletBuilder servlet = new MarquezGraphqlServletBuilder();
168179
this.graphqlServlet = servlet.getServlet(new GraphqlSchemaBuilder(jdbi));

api/src/main/java/marquez/api/BaseResource.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
import marquez.service.DatasetFieldService;
2929
import marquez.service.DatasetService;
3030
import marquez.service.DatasetVersionService;
31+
import marquez.service.EventService;
3132
import marquez.service.JobService;
3233
import marquez.service.LineageService;
3334
import marquez.service.NamespaceService;
@@ -50,6 +51,7 @@ public class BaseResource {
5051
protected DatasetVersionService datasetVersionService;
5152
protected DatasetFieldService datasetFieldService;
5253
protected LineageService lineageService;
54+
protected EventService eventService;
5355

5456
public BaseResource(ServiceFactory serviceFactory) {
5557
this.serviceFactory = serviceFactory;
@@ -63,6 +65,7 @@ public BaseResource(ServiceFactory serviceFactory) {
6365
this.datasetVersionService = serviceFactory.getDatasetVersionService();
6466
this.datasetFieldService = serviceFactory.getDatasetFieldService();
6567
this.lineageService = serviceFactory.getLineageService();
68+
this.eventService = serviceFactory.getEventService();
6669
}
6770

6871
void throwIfNotExists(@NonNull NamespaceName namespaceName) {
Lines changed: 62 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,62 @@
1+
package marquez.api;
2+
3+
import static javax.ws.rs.core.MediaType.APPLICATION_JSON;
4+
5+
import com.codahale.metrics.annotation.ExceptionMetered;
6+
import com.codahale.metrics.annotation.ResponseMetered;
7+
import com.codahale.metrics.annotation.Timed;
8+
import com.fasterxml.jackson.annotation.JsonProperty;
9+
import com.fasterxml.jackson.databind.JsonNode;
10+
import java.util.List;
11+
import javax.validation.constraints.Min;
12+
import javax.ws.rs.DefaultValue;
13+
import javax.ws.rs.GET;
14+
import javax.ws.rs.Path;
15+
import javax.ws.rs.PathParam;
16+
import javax.ws.rs.Produces;
17+
import javax.ws.rs.QueryParam;
18+
import javax.ws.rs.core.Response;
19+
import lombok.NonNull;
20+
import lombok.Value;
21+
import marquez.service.ServiceFactory;
22+
23+
@Path("/api/v1")
24+
public class EventResource extends BaseResource {
25+
public EventResource(@NonNull final ServiceFactory serviceFactory) {
26+
super(serviceFactory);
27+
}
28+
29+
@Timed
30+
@ResponseMetered
31+
@ExceptionMetered
32+
@GET
33+
@Path("/events")
34+
@Produces(APPLICATION_JSON)
35+
public Response get(
36+
@QueryParam("limit") @DefaultValue("100") @Min(value = 0) int limit,
37+
@QueryParam("offset") @DefaultValue("0") @Min(value = 0) int offset) {
38+
final List<JsonNode> events = eventService.getAll(limit, offset);
39+
return Response.ok(new Events(events)).build();
40+
}
41+
42+
@Timed
43+
@ResponseMetered
44+
@ExceptionMetered
45+
@GET
46+
@Path("/events/{namespace}")
47+
@Produces(APPLICATION_JSON)
48+
public Response getByNamespace(
49+
@PathParam("namespace") String namespace,
50+
@QueryParam("limit") @DefaultValue("100") @Min(value = 0) int limit,
51+
@QueryParam("offset") @DefaultValue("0") @Min(value = 0) int offset) {
52+
final List<JsonNode> event = eventService.getByNamespace(namespace, limit, offset);
53+
return Response.ok(new Events(event)).build();
54+
}
55+
56+
@Value
57+
static class Events {
58+
@NonNull
59+
@JsonProperty("events")
60+
List<JsonNode> value;
61+
}
62+
}

api/src/main/java/marquez/db/BaseDao.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,4 +50,7 @@ public interface BaseDao extends SqlObject {
5050

5151
@CreateSqlObject
5252
OpenLineageDao createOpenLineageDao();
53+
54+
@CreateSqlObject
55+
EventDao createEventDao();
5356
}

api/src/main/java/marquez/db/Columns.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -126,6 +126,9 @@ private Columns() {}
126126
public static final String RUN_UUID = "run_uuid";
127127
public static final String STATE = "state";
128128

129+
/* LINEAGE EVENT ROW COLUMNS */
130+
public static final String EVENT = "event";
131+
129132
public static UUID uuidOrNull(final ResultSet results, final String column) throws SQLException {
130133
if (results.getObject(column) == null) {
131134
return null;
Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,37 @@
1+
package marquez.db;
2+
3+
import com.fasterxml.jackson.databind.JsonNode;
4+
import java.util.List;
5+
import marquez.db.mappers.RawLineageEventMapper;
6+
import org.jdbi.v3.sqlobject.config.RegisterRowMapper;
7+
import org.jdbi.v3.sqlobject.customizer.Bind;
8+
import org.jdbi.v3.sqlobject.statement.SqlQuery;
9+
10+
@RegisterRowMapper(RawLineageEventMapper.class)
11+
public interface EventDao extends BaseDao {
12+
13+
@SqlQuery(
14+
"""
15+
SELECT event
16+
FROM lineage_events
17+
ORDER BY event_time DESC
18+
LIMIT :limit
19+
OFFSET :offset""")
20+
List<JsonNode> getAll(int limit, int offset);
21+
22+
/**
23+
* This is a "hack" to get inputs/outputs namespace from jsonb column: <a
24+
* href="https://github.com/jdbi/jdbi/issues/1510#issuecomment-485423083">explanation</a>
25+
*/
26+
@SqlQuery(
27+
"""
28+
SELECT le.event
29+
FROM lineage_events le, jsonb_array_elements(coalesce(le.event -> 'inputs', '[]'::jsonb) || coalesce(le.event -> 'outputs', '[]'::jsonb)) AS ds
30+
WHERE le.job_namespace = :namespace
31+
OR ds ->> 'namespace' = :namespace
32+
ORDER BY event_time DESC
33+
LIMIT :limit
34+
OFFSET :offset""")
35+
List<JsonNode> getByNamespace(
36+
@Bind("namespace") String namespace, @Bind("limit") int limit, @Bind("offset") int offset);
37+
}
Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,30 @@
1+
package marquez.db.mappers;
2+
3+
import static marquez.db.Columns.stringOrThrow;
4+
5+
import com.fasterxml.jackson.core.JsonProcessingException;
6+
import com.fasterxml.jackson.databind.JsonNode;
7+
import com.fasterxml.jackson.databind.ObjectMapper;
8+
import java.sql.ResultSet;
9+
import java.sql.SQLException;
10+
import lombok.extern.slf4j.Slf4j;
11+
import marquez.common.Utils;
12+
import marquez.db.Columns;
13+
import org.jdbi.v3.core.mapper.RowMapper;
14+
import org.jdbi.v3.core.statement.StatementContext;
15+
16+
@Slf4j
17+
public class RawLineageEventMapper implements RowMapper<JsonNode> {
18+
@Override
19+
public JsonNode map(ResultSet rs, StatementContext ctx) throws SQLException {
20+
String rawEvent = stringOrThrow(rs, Columns.EVENT);
21+
22+
try {
23+
ObjectMapper mapper = Utils.getMapper();
24+
return mapper.readTree(rawEvent);
25+
} catch (JsonProcessingException e) {
26+
log.error("Failed to process json", e);
27+
}
28+
return null;
29+
}
30+
}

api/src/main/java/marquez/service/DelegatingDaos.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@
1010
import marquez.db.DatasetDao;
1111
import marquez.db.DatasetFieldDao;
1212
import marquez.db.DatasetVersionDao;
13+
import marquez.db.EventDao;
1314
import marquez.db.JobContextDao;
1415
import marquez.db.JobDao;
1516
import marquez.db.JobVersionDao;
@@ -98,4 +99,9 @@ public static class DelegatingTagDao implements TagDao {
9899
public static class DelegatingLineageDao implements LineageDao {
99100
@Delegate private final LineageDao delegate;
100101
}
102+
103+
@RequiredArgsConstructor
104+
public static class DelegatingEventDao implements EventDao {
105+
@Delegate private final EventDao delegate;
106+
}
101107
}
Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,10 @@
1+
package marquez.service;
2+
3+
import lombok.NonNull;
4+
import marquez.db.BaseDao;
5+
6+
public class EventService extends DelegatingDaos.DelegatingEventDao {
7+
public EventService(@NonNull BaseDao baseDao) {
8+
super(baseDao.createEventDao());
9+
}
10+
}

0 commit comments

Comments
 (0)