Skip to content

BigQuery: 'test_to_dataframe_w_bqstorage_nonempty' unit test flakes #8175

Description

@tseaver

From this failed Kokoro job:

____________ TestRowIterator.test_to_dataframe_w_bqstorage_nonempty ____________
self = <tests.unit.test_table.TestRowIterator testMethod=test_to_dataframe_w_bqstorage_nonempty>
    @unittest.skipIf(pandas is None, "Requires `pandas`")
    @unittest.skipIf(
        bigquery_storage_v1beta1 is None, "Requires `google-cloud-bigquery-storage`"
    )
    def test_to_dataframe_w_bqstorage_nonempty(self):
        from google.cloud.bigquery import schema
        from google.cloud.bigquery import table as mut
        from google.cloud.bigquery_storage_v1beta1 import reader
        # Speed up testing.
        mut._PROGRESS_INTERVAL = 0.01
        bqstorage_client = mock.create_autospec(
            bigquery_storage_v1beta1.BigQueryStorageClient
        )
        streams = [
            # Use two streams we want to check frames are read from each stream.
            {"name": "/projects/proj/dataset/dset/tables/tbl/streams/1234"},
            {"name": "/projects/proj/dataset/dset/tables/tbl/streams/5678"},
        ]
        session = bigquery_storage_v1beta1.types.ReadSession(streams=streams)
        session.avro_schema.schema = json.dumps(
            {
                "fields": [
                    {"name": "colA"},
                    # Not alphabetical to test column order.
                    {"name": "colC"},
                    {"name": "colB"},
                ]
            }
        )
        bqstorage_client.create_read_session.return_value = session
        mock_rowstream = mock.create_autospec(reader.ReadRowsStream)
        bqstorage_client.read_rows.return_value = mock_rowstream
        mock_rows = mock.create_autospec(reader.ReadRowsIterable)
        mock_rowstream.rows.return_value = mock_rows
        page_items = [
            {"colA": 1, "colB": "abc", "colC": 2.0},
            {"colA": -1, "colB": "def", "colC": 4.0},
        ]
        def blocking_to_dataframe(*args, **kwargs):
            # Sleep for longer than the waiting interval so that we know we're
            # only reading one page per loop at most.
            time.sleep(2 * mut._PROGRESS_INTERVAL)
            return pandas.DataFrame(page_items, columns=["colA", "colB", "colC"])
        mock_page = mock.create_autospec(reader.ReadRowsPage)
        mock_page.to_dataframe.side_effect = blocking_to_dataframe
        mock_pages = (mock_page, mock_page, mock_page)
        type(mock_rows).pages = mock.PropertyMock(return_value=mock_pages)
        # Test that full queue errors are ignored.
        mock_queue = mock.create_autospec(mut._NoopProgressBarQueue)
        mock_queue().put_nowait.side_effect = queue.Full
        schema = [
            schema.SchemaField("colA", "IGNORED"),
            schema.SchemaField("colC", "IGNORED"),
            schema.SchemaField("colB", "IGNORED"),
        ]
        row_iterator = mut.RowIterator(
            _mock_client(),
            None,  # api_request: ignored
            None,  # path: ignored
            schema,
            table=mut.TableReference.from_string("proj.dset.tbl"),
            selected_fields=schema,
        )
        with mock.patch.object(mut, "_NoopProgressBarQueue", mock_queue), mock.patch(
            "concurrent.futures.wait", wraps=concurrent.futures.wait
        ) as mock_wait:
            got = row_iterator.to_dataframe(bqstorage_client=bqstorage_client)
        # Are the columns in the expected order?
        column_names = ["colA", "colC", "colB"]
        self.assertEqual(list(got), column_names)
        # Have expected number of rows?
        total_pages = len(streams) * len(mock_pages)
        total_rows = len(page_items) * total_pages
        self.assertEqual(len(got.index), total_rows)
        # Make sure that this test looped through multiple progress intervals.
        self.assertGreaterEqual(mock_wait.call_count, 2)
        # Make sure that this test pushed to the progress queue.
>       self.assertEqual(mock_queue().put_nowait.call_count, total_pages)
E AssertionError: 5 != 6

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

api: bigqueryIssues related to the BigQuery API.flakytestingtype: processA process-related concern. May include testing, release, or the like.

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions