7777 spark = SparkSession.builder.appName("test-hbase").getOrCreate()
7878
7979 df = spark.createDataFrame(
80- [("row1", "Hello, Stackable!")],
80+ [
81+ ("row1", "Hello, Stackable!"),
82+ ("row2", "Hello, HBase!"),
83+ ("row3", "Hello, Spark!"),
84+ ],
8185 "key: string, value: string"
8286 )
8387
@@ -101,3 +105,53 @@ data:
101105 .option('catalog', catalog)\
102106 .option('newtable', '5')\
103107 .save()
108+
109+ # Read the whole table back and register it as a temporary view so we can
110+ # run SQL queries against it.
111+ #
112+ # hbase.spark.pushdown.columnfilter=false disables server-side predicate
113+ # pushdown: the connector would otherwise ship a SparkSQLPushDownFilter to
114+ # the RegionServers, which requires the hbase-spark connector classes on
115+ # the RegionServer classpath (not present in the Stackable HBase image).
116+ # With pushdown off, WHERE clauses are evaluated Spark-side over a full
117+ # scan, which is all this test needs.
118+ read_df = spark\
119+ .read\
120+ .format("org.apache.hadoop.hbase.spark")\
121+ .option('catalog', catalog)\
122+ .option('hbase.spark.pushdown.columnfilter', 'false')\
123+ .load()
124+ read_df.createOrReplaceTempView("test_hbase")
125+
126+ # Full scan: all rows that were written must be readable again.
127+ all_rows = spark.sql("SELECT key, value FROM test_hbase").collect()
128+ actual = {row["key"]: row["value"] for row in all_rows}
129+ expected = {
130+ "row1": "Hello, Stackable!",
131+ "row2": "Hello, HBase!",
132+ "row3": "Hello, Spark!",
133+ }
134+ assert actual == expected, f"full scan mismatch: {actual} != {expected}"
135+
136+ # Row-key point lookup.
137+ point = spark.sql("SELECT value FROM test_hbase WHERE key = 'row2'").collect()
138+ assert len(point) == 1, f"point lookup returned {len(point)} rows, expected 1"
139+ assert point[0]["value"] == "Hello, HBase!", f"unexpected value: {point[0]['value']}"
140+
141+ # Range scan on the row key.
142+ range_rows = spark.sql(
143+ "SELECT key FROM test_hbase WHERE key >= 'row2' ORDER BY key"
144+ ).collect()
145+ assert [r["key"] for r in range_rows] == ["row2", "row3"], \
146+ f"unexpected range scan result: {[r['key'] for r in range_rows]}"
147+
148+ # Filter on a non-rowkey column value.
149+ like_rows = spark.sql(
150+ "SELECT key FROM test_hbase WHERE value LIKE 'Hello, S%' ORDER BY key"
151+ ).collect()
152+ assert [r["key"] for r in like_rows] == ["row1", "row3"], \
153+ f"unexpected value filter result: {[r['key'] for r in like_rows]}"
154+
155+ # Aggregation over the scanned rows.
156+ count = spark.sql("SELECT COUNT(*) AS c FROM test_hbase").collect()[0]["c"]
157+ assert count == 3, f"expected 3 rows, got {count}"
0 commit comments